1use std::{
19 future::Future,
20 sync::Arc,
21 time::{Duration, Instant},
22};
23
24use anyhow::Context;
25use async_trait::async_trait;
26use jiff::Timestamp;
27use nautilus_common::{
28 cache::InstrumentLookupError,
29 clients::ExecutionClient,
30 live::runner::get_exec_event_sender,
31 messages::execution::{
32 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
33 GenerateOrderStatusReport, GenerateOrderStatusReports, GeneratePositionStatusReports,
34 ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
35 },
36};
37use nautilus_core::{
38 AtomicMap, Params, UnixNanos,
39 time::{AtomicTime, get_atomic_clock_realtime},
40};
41use nautilus_live::{
42 ExecutionClientCore, ExecutionEventEmitter, SocketControl, execution::failure::CommandFailure,
43 task::TaskGroup,
44};
45use nautilus_model::{
46 accounts::AccountAny,
47 enums::{AccountType, OmsType, OrderStatus, OrderType},
48 identifiers::{
49 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, Venue, VenueOrderId,
50 },
51 instruments::{Instrument, InstrumentAny},
52 orders::{Order, OrderAny},
53 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
54 types::{AccountBalance, MarginBalance, Quantity},
55};
56use rust_decimal::Decimal;
57use tokio_util::sync::CancellationToken;
58
59use super::{
60 command_failure_from_cancel_error, command_failure_from_futures_batch_error,
61 command_failure_from_futures_batch_item, command_failure_from_modify_error,
62 command_failure_from_submit_error,
63};
64use crate::{
65 common::{
66 consts::KRAKEN_VENUE,
67 credential::KrakenCredential,
68 enums::{KrakenApiResult, KrakenSendStatus},
69 parse::truncate_cl_ord_id,
70 },
71 config::KrakenExecutionClientConfig,
72 http::{
73 KrakenFuturesHttpClient,
74 futures::{
75 client::KRAKEN_FUTURES_DEFAULT_RATE_LIMIT_PER_SECOND, models::FuturesBatchCancelStatus,
76 query::KrakenFuturesBatchCancelItem,
77 },
78 },
79 websocket::{
80 dispatch::{self, OrderIdentity, WsDispatchState},
81 futures::{client::KrakenFuturesWebSocketClient, messages::KrakenFuturesWsMessage},
82 },
83};
84
85const FUTURES_BATCH_CANCEL_LIMIT: usize = 50;
86
87#[allow(dead_code)]
92#[derive(Debug)]
93pub struct KrakenFuturesExecutionClient {
94 core: ExecutionClientCore,
95 clock: &'static AtomicTime,
96 config: KrakenExecutionClientConfig,
97 emitter: ExecutionEventEmitter,
98 http: KrakenFuturesHttpClient,
99 ws: KrakenFuturesWebSocketClient,
100 cancellation_token: CancellationToken,
101 session_tasks: TaskGroup,
102 pending_tasks: TaskGroup,
103 instruments: Arc<AtomicMap<InstrumentId, InstrumentAny>>,
104 truncated_id_map: Arc<AtomicMap<String, ClientOrderId>>,
105 order_instrument_map: Arc<AtomicMap<String, InstrumentId>>,
106 venue_client_map: Arc<AtomicMap<String, ClientOrderId>>,
107 venue_order_qty: Arc<AtomicMap<String, Quantity>>,
108 ws_dispatch_state: Arc<WsDispatchState>,
109}
110
111impl KrakenFuturesExecutionClient {
112 pub fn new(
114 core: ExecutionClientCore,
115 config: KrakenExecutionClientConfig,
116 ) -> anyhow::Result<Self> {
117 let clock = get_atomic_clock_realtime();
118 let emitter = ExecutionEventEmitter::new(
119 clock,
120 core.trader_id,
121 core.account_id,
122 AccountType::Margin,
123 None,
124 );
125
126 let session_tasks = TaskGroup::new();
127 let cancellation_token = session_tasks.cancellation_token();
128 let pending_tasks = TaskGroup::new();
129
130 let http = KrakenFuturesHttpClient::with_credentials(
131 config.api_key.clone(),
132 config.api_secret.clone(),
133 config.environment,
134 config.base_url.clone(),
135 config.timeout_secs,
136 None,
137 None,
138 None,
139 config.proxy_url.clone(),
140 config
141 .max_requests_per_second
142 .unwrap_or(KRAKEN_FUTURES_DEFAULT_RATE_LIMIT_PER_SECOND),
143 )?;
144
145 let credential = KrakenCredential::new(config.api_key.clone(), config.api_secret.clone());
146 let ws = KrakenFuturesWebSocketClient::with_credentials(
147 config.ws_url(),
148 config.heartbeat_interval_secs,
149 Some(credential),
150 config.auth_timeout_secs,
151 config.transport_backend,
152 config.proxy_url.clone(),
153 )
154 .with_socket_control(SocketControl::new(
155 core.client_id,
156 Some(*KRAKEN_VENUE),
157 "kraken-futures-user-streams",
158 ));
159
160 Ok(Self {
161 core,
162 clock,
163 config,
164 emitter,
165 http,
166 ws,
167 cancellation_token,
168 session_tasks,
169 pending_tasks,
170 instruments: Arc::new(AtomicMap::new()),
171 truncated_id_map: Arc::new(AtomicMap::new()),
172 order_instrument_map: Arc::new(AtomicMap::new()),
173 venue_client_map: Arc::new(AtomicMap::new()),
174 venue_order_qty: Arc::new(AtomicMap::new()),
175 ws_dispatch_state: Arc::new(WsDispatchState::new()),
176 })
177 }
178
179 fn register_order_identity(&self, order: &OrderAny) {
180 self.ws_dispatch_state.register_identity(
181 order.client_order_id(),
182 OrderIdentity {
183 strategy_id: order.strategy_id(),
184 instrument_id: order.instrument_id(),
185 order_side: order.order_side(),
186 order_type: order.order_type(),
187 quantity: order.quantity(),
188 },
189 );
190 }
191
192 #[must_use]
194 pub fn clock(&self) -> &'static AtomicTime {
195 self.clock
196 }
197
198 #[must_use]
200 pub fn emitter(&self) -> &ExecutionEventEmitter {
201 &self.emitter
202 }
203
204 fn spawn_task<F>(&self, description: &'static str, fut: F)
205 where
206 F: Future<Output = anyhow::Result<()>> + Send + 'static,
207 {
208 let future = async move {
209 if let Err(e) = fut.await {
210 log::warn!("{description} failed: {e:?}");
211 }
212 };
213
214 if let Err(e) = self.pending_tasks.spawn(future) {
215 log::warn!("Skipping Kraken Futures {description} after shutdown began: {e}");
216 }
217 }
218
219 async fn finish_tasks(&self) -> anyhow::Result<()> {
220 let (session_result, pending_result) = tokio::join!(
221 self.session_tasks
222 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
223 self.pending_tasks
224 .finish_shutdown(Duration::from_secs(2), Duration::from_secs(2)),
225 );
226 session_result.context("failed to finish Kraken Futures execution session tasks")?;
227 pending_result.context("failed to finish Kraken Futures execution command tasks")?;
228 Ok(())
229 }
230
231 async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
232 if !self.session_tasks.is_open() || !self.pending_tasks.is_open() {
233 self.session_tasks.begin_shutdown();
234 self.pending_tasks.begin_shutdown();
235 self.finish_tasks().await?;
236 self.session_tasks
237 .start_generation()
238 .context("failed to start Kraken Futures execution session task generation")?;
239 self.pending_tasks
240 .start_generation()
241 .context("failed to start Kraken Futures execution command task generation")?;
242 self.cancellation_token = self.session_tasks.cancellation_token();
243 }
244 Ok(())
245 }
246
247 async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
248 self.http.cancel_all_requests();
249 self.pending_tasks.begin_shutdown();
250 let ws_result = self.ws.close().await;
251 self.session_tasks.begin_shutdown();
252 let tasks_result = self.finish_tasks().await;
253 self.core.set_disconnected();
254 tasks_result?;
255 Ok(ws_result?)
256 }
257
258 fn submit_single_order(&self, order: &OrderAny, task_name: &'static str) {
259 if order.is_closed() {
260 log::warn!(
261 "Cannot submit closed order: client_order_id={}",
262 order.client_order_id()
263 );
264 return;
265 }
266
267 let account_id = self.core.account_id;
268 let client_order_id = order.client_order_id();
269 let strategy_id = order.strategy_id();
270 let instrument_id = order.instrument_id();
271 let order_side = order.order_side();
272 let order_type = order.order_type();
273 let quantity = order.quantity();
274 let time_in_force = order.time_in_force();
275 let price = order.price();
276 let trigger_price = order.trigger_price();
277 let trigger_type = order.trigger_type();
278 let is_reduce_only = order.is_reduce_only();
279 let is_post_only = order.is_post_only();
280
281 log::debug!("OrderSubmitted: client_order_id={client_order_id}");
282 self.register_order_identity(order);
283 self.emitter.emit_order_submitted(order);
284
285 let kraken_cl_ord_id = truncate_cl_ord_id(&client_order_id);
286
287 if kraken_cl_ord_id != client_order_id.as_str() {
288 self.truncated_id_map
289 .insert(kraken_cl_ord_id, client_order_id);
290 }
291
292 let http = self.http.clone();
293 let emitter = self.emitter.clone();
294 let clock = self.clock;
295 let dispatch_state = self.ws_dispatch_state.clone();
296
297 self.spawn_task(task_name, async move {
298 let result = http
299 .submit_order(
300 account_id,
301 instrument_id,
302 client_order_id,
303 order_side,
304 order_type,
305 quantity,
306 time_in_force,
307 price,
308 trigger_price,
309 trigger_type,
310 is_reduce_only,
311 is_post_only,
312 )
313 .await;
314
315 match result {
316 Ok(_) => {}
317 Err(e) => match command_failure_from_submit_error(&e) {
318 CommandFailure::Ambiguous(reason) => {
319 log::warn!(
320 "{task_name} outcome is ambiguous for client_order_id={client_order_id}: {reason}"
321 );
322 }
323 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
324 let ts_event = clock.get_time_ns();
325 let error_msg = format!("{task_name} error: {reason}");
326 let due_post_only = error_msg.contains("POST_ONLY_REJECTED");
327 dispatch_state.cleanup_terminal(&client_order_id);
328 emitter.emit_order_rejected_event(
329 strategy_id,
330 instrument_id,
331 client_order_id,
332 &error_msg,
333 ts_event,
334 due_post_only,
335 );
336 }
337 },
338 }
339 Ok(())
340 });
341 }
342
343 fn cancel_single_order(&self, cmd: &CancelOrder) {
344 let account_id = self.core.account_id;
345 let client_order_id = cmd.client_order_id;
346 let venue_order_id = cmd.venue_order_id;
347 let strategy_id = cmd.strategy_id;
348 let instrument_id = cmd.instrument_id;
349
350 log::debug!(
351 "Canceling order: venue_order_id={venue_order_id:?}, client_order_id={client_order_id}"
352 );
353
354 let http = self.http.clone();
355 let emitter = self.emitter.clone();
356 let clock = self.clock;
357
358 self.spawn_task("cancel_order", async move {
359 if let Err(failure) = cancel_order_for_futures(
360 &http,
361 account_id,
362 instrument_id,
363 Some(client_order_id),
364 venue_order_id,
365 )
366 .await
367 {
368 handle_cancel_failure(
369 &emitter,
370 clock,
371 strategy_id,
372 instrument_id,
373 client_order_id,
374 venue_order_id,
375 failure,
376 );
377 }
378 Ok(())
379 });
380 }
381
382 fn spawn_message_handler(&mut self) -> anyhow::Result<()> {
383 let mut rx = self
384 .ws
385 .take_output_rx()
386 .context("Failed to take futures WebSocket output receiver")?;
387 let emitter = self.emitter.clone();
388 let instruments = self.instruments.clone();
389 let truncated_id_map = self.truncated_id_map.clone();
390 let order_instrument_map = self.order_instrument_map.clone();
391 let venue_client_map = self.venue_client_map.clone();
392 let venue_order_qty = self.venue_order_qty.clone();
393 let dispatch_state = self.ws_dispatch_state.clone();
394 let account_id = self.core.account_id;
395 let clock = self.clock;
396 let cancellation_token = self.cancellation_token.clone();
397
398 let future = async move {
399 loop {
400 tokio::select! {
401 () = cancellation_token.cancelled() => {
402 log::debug!("Futures execution message handler cancelled");
403 break;
404 }
405 msg = rx.recv() => {
406 match msg {
407 Some(ws_msg) => {
408 Self::handle_ws_message(
409 ws_msg,
410 &emitter,
411 &dispatch_state,
412 &instruments,
413 &truncated_id_map,
414 &order_instrument_map,
415 &venue_client_map,
416 &venue_order_qty,
417 account_id,
418 clock,
419 );
420 }
421 None => {
422 log::debug!("Futures execution WebSocket stream ended");
423 break;
424 }
425 }
426 }
427 }
428 }
429 };
430
431 self.session_tasks
432 .spawn(future)
433 .context("failed to register Kraken Futures execution stream task")
434 }
435
436 #[expect(clippy::too_many_arguments)]
437 fn handle_ws_message(
438 msg: KrakenFuturesWsMessage,
439 emitter: &ExecutionEventEmitter,
440 dispatch_state: &Arc<WsDispatchState>,
441 instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
442 truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
443 order_instrument_map: &Arc<AtomicMap<String, InstrumentId>>,
444 venue_client_map: &Arc<AtomicMap<String, ClientOrderId>>,
445 venue_order_qty: &Arc<AtomicMap<String, Quantity>>,
446 account_id: AccountId,
447 clock: &'static AtomicTime,
448 ) {
449 let ts_init = clock.get_time_ns();
450
451 match msg {
452 KrakenFuturesWsMessage::OpenOrdersDelta(delta) => {
453 dispatch::futures::open_orders_delta(
454 &delta,
455 dispatch_state,
456 emitter,
457 instruments,
458 truncated_id_map,
459 order_instrument_map,
460 venue_client_map,
461 venue_order_qty,
462 account_id,
463 ts_init,
464 );
465 }
466 KrakenFuturesWsMessage::OpenOrdersCancel(cancel) => {
467 dispatch::futures::open_orders_cancel(
468 &cancel,
469 dispatch_state,
470 emitter,
471 truncated_id_map,
472 order_instrument_map,
473 venue_client_map,
474 venue_order_qty,
475 account_id,
476 ts_init,
477 );
478 }
479 KrakenFuturesWsMessage::FillsDelta(fills_delta) => {
480 dispatch::futures::fills_delta(
481 &fills_delta,
482 dispatch_state,
483 emitter,
484 instruments,
485 truncated_id_map,
486 venue_client_map,
487 account_id,
488 ts_init,
489 );
490 }
491 KrakenFuturesWsMessage::Challenge(challenge) => {
492 log::debug!("Received challenge: length={}", challenge.len());
493 }
494 KrakenFuturesWsMessage::Reconnected => {
495 log::info!("Futures execution WebSocket reconnected");
496 }
497 KrakenFuturesWsMessage::Ticker(_)
498 | KrakenFuturesWsMessage::Trade(_)
499 | KrakenFuturesWsMessage::BookSnapshot(_)
500 | KrakenFuturesWsMessage::BookDelta(_) => {}
501 }
502 }
503
504 async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
505 let account_id = self.core.account_id;
506
507 if self.core.cache().account(&account_id).is_some() {
508 log::info!("Account {account_id} registered");
509 return Ok(());
510 }
511
512 let start = Instant::now();
513 let timeout = Duration::from_secs_f64(timeout_secs);
514 let interval = Duration::from_millis(10);
515
516 loop {
517 tokio::time::sleep(interval).await;
518
519 if self.core.cache().account(&account_id).is_some() {
520 log::info!("Account {account_id} registered");
521 return Ok(());
522 }
523
524 if start.elapsed() >= timeout {
525 anyhow::bail!(
526 "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
527 );
528 }
529 }
530 }
531
532 fn modify_single_order(&self, cmd: &ModifyOrder) {
533 let client_order_id = cmd.client_order_id;
534 let venue_order_id = cmd.venue_order_id;
535 let strategy_id = cmd.strategy_id;
536 let instrument_id = cmd.instrument_id;
537 let quantity = cmd.quantity;
538 let price = cmd.price;
539
540 log::debug!(
541 "Modifying order: venue_order_id={venue_order_id:?}, client_order_id={client_order_id}"
542 );
543
544 let http = self.http.clone();
545 let emitter = self.emitter.clone();
546 let clock = self.clock;
547
548 self.spawn_task("modify_order", async move {
549 match http
550 .modify_order(
551 instrument_id,
552 Some(client_order_id),
553 venue_order_id,
554 quantity,
555 price,
556 None,
557 )
558 .await
559 {
560 Ok(_) => {}
561 Err(e) => match command_failure_from_modify_error(&e) {
562 CommandFailure::Ambiguous(reason) => {
563 log::warn!(
564 "modify_order outcome is ambiguous for client_order_id={client_order_id}: {reason}"
565 );
566 }
567 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
568 let ts_event = clock.get_time_ns();
569 emitter.emit_order_modify_rejected_event(
570 strategy_id,
571 instrument_id,
572 client_order_id,
573 venue_order_id,
574 &format!("modify-order error: {reason}"),
575 ts_event,
576 );
577 }
578 },
579 }
580 Ok(())
581 });
582 }
583}
584
585#[async_trait(?Send)]
586impl ExecutionClient for KrakenFuturesExecutionClient {
587 fn is_connected(&self) -> bool {
588 self.core.is_connected()
589 }
590
591 fn client_id(&self) -> ClientId {
592 self.core.client_id
593 }
594
595 fn account_id(&self) -> AccountId {
596 self.core.account_id
597 }
598
599 fn venue(&self) -> Venue {
600 *KRAKEN_VENUE
601 }
602
603 fn oms_type(&self) -> OmsType {
604 self.core.oms_type
605 }
606
607 fn get_account(&self) -> Option<AccountAny> {
608 self.core.cache().account_owned(&self.core.account_id)
609 }
610
611 fn generate_account_state(
612 &self,
613 balances: Vec<AccountBalance>,
614 margins: Vec<MarginBalance>,
615 reported: bool,
616 ts_event: UnixNanos,
617 info: Option<Params>,
618 ) -> anyhow::Result<()> {
619 self.emitter
620 .emit_account_state(balances, margins, reported, ts_event, info);
621 Ok(())
622 }
623
624 fn start(&mut self) -> anyhow::Result<()> {
625 if self.core.is_started() {
626 return Ok(());
627 }
628
629 self.emitter.set_sender(get_exec_event_sender());
630 self.core.set_started();
631
632 log::info!(
633 "Started: client_id={}, account_id={}, product_type=Futures, environment={:?}",
634 self.core.client_id,
635 self.core.account_id,
636 self.config.environment
637 );
638 Ok(())
639 }
640
641 fn stop(&mut self) -> anyhow::Result<()> {
642 if self.core.is_stopped() {
643 return Ok(());
644 }
645
646 self.http.cancel_all_requests();
647 self.session_tasks.begin_shutdown();
648 self.pending_tasks.begin_shutdown();
649 self.ws.begin_shutdown();
650 self.core.set_stopped();
651 self.core.set_disconnected();
652 log::info!("Stopped: client_id={}", self.core.client_id);
653 Ok(())
654 }
655
656 async fn connect(&mut self) -> anyhow::Result<()> {
657 if self.core.is_connected() && self.session_tasks.is_open() && self.pending_tasks.is_open()
658 {
659 return Ok(());
660 }
661
662 self.http.reset_cancellation_token();
663 self.prepare_task_groups().await?;
664
665 if !self.core.instruments_initialized() {
666 let instruments = self
667 .http
668 .request_instruments()
669 .await
670 .context("Failed to load Kraken futures instruments")?;
671 log::debug!("Loaded {} Futures instruments", instruments.len());
672 self.http.cache_instruments(&instruments);
673 self.core.set_instruments_initialized();
674 }
675
676 self.instruments.rcu(|m| {
677 for instrument in self.http.instruments_cache.load().values() {
678 m.insert(instrument.id(), instrument.clone());
679 }
680 });
681
682 let session_result = async {
683 self.ws
684 .connect()
685 .await
686 .context("Failed to connect futures WebSocket")?;
687 self.ws
688 .wait_until_active(10.0)
689 .await
690 .context("Futures WebSocket failed to become active")?;
691
692 self.ws
693 .authenticate()
694 .await
695 .context("Failed to authenticate futures WebSocket")?;
696
697 let account_state = self
698 .http
699 .request_account_state(self.core.account_id)
700 .await
701 .context("Failed to request Kraken futures account state")?;
702
703 if !account_state.balances.is_empty() {
704 log::debug!(
705 "Received account state with {} balance(s)",
706 account_state.balances.len()
707 );
708 }
709 self.emitter.send_account_state(account_state);
710 self.await_account_registered(30.0).await?;
711
712 self.spawn_message_handler()?;
713
714 self.ws
715 .subscribe_executions()
716 .await
717 .context("Failed to subscribe to executions")?;
718
719 log::debug!("Futures WebSocket authenticated and subscribed to executions");
720
721 Ok::<(), anyhow::Error>(())
722 }
723 .await;
724
725 if let Err(e) = session_result {
726 if let Err(teardown_error) = self.teardown_partial_connect().await {
727 return Err(e.context(format!(
728 "Kraken Futures execution startup teardown failed: {teardown_error}"
729 )));
730 }
731 return Err(e);
732 }
733
734 self.core.set_connected();
735 log::info!("Connected: client_id={}", self.core.client_id);
736 Ok(())
737 }
738
739 async fn disconnect(&mut self) -> anyhow::Result<()> {
740 self.teardown_partial_connect().await?;
741 log::info!("Disconnected: client_id={}", self.core.client_id);
742 Ok(())
743 }
744
745 async fn generate_order_status_report(
746 &self,
747 cmd: &GenerateOrderStatusReport,
748 ) -> anyhow::Result<Option<OrderStatusReport>> {
749 log::debug!(
750 "Generating order status report: venue_order_id={:?}, client_order_id={:?}",
751 cmd.venue_order_id,
752 cmd.client_order_id
753 );
754
755 let account_id = self.core.account_id;
756 let reports = self
757 .http
758 .request_order_status_reports(account_id, None, None, None, false)
759 .await?;
760
761 let matched = reports.into_iter().find(|r| {
764 cmd.venue_order_id
765 .is_some_and(|id| r.venue_order_id.as_str() == id.as_str())
766 || cmd.client_order_id.is_some_and(|id| {
767 r.client_order_id
768 .as_ref()
769 .is_some_and(|r_id| r_id.as_str() == truncate_cl_ord_id(&id))
770 })
771 });
772
773 if matched.is_some() {
774 return Ok(matched);
775 }
776
777 let Some(order) = self.get_cached_order_for_status_command(cmd) else {
778 return Ok(None);
779 };
780
781 let now = Timestamp::now();
782 let start = now - Duration::from_secs(5 * 60);
783 let fills = self
784 .http
785 .request_fill_reports(
786 account_id,
787 Some(order.instrument_id()),
788 Some(start),
789 Some(now),
790 )
791 .await?;
792
793 Ok(synthesize_filled_order_status_report(cmd, &order, &fills))
794 }
795
796 async fn generate_order_status_reports(
797 &self,
798 cmd: &GenerateOrderStatusReports,
799 ) -> anyhow::Result<Vec<OrderStatusReport>> {
800 log::debug!(
801 "Generating order status reports: instrument_id={:?}, open_only={}",
802 cmd.instrument_id,
803 cmd.open_only
804 );
805
806 let account_id = self.core.account_id;
807 let start = cmd.start.map(Timestamp::from);
808 let end = cmd.end.map(Timestamp::from);
809 self.http
810 .request_order_status_reports(account_id, cmd.instrument_id, start, end, cmd.open_only)
811 .await
812 }
813
814 async fn generate_fill_reports(
815 &self,
816 cmd: GenerateFillReports,
817 ) -> anyhow::Result<Vec<FillReport>> {
818 log::debug!(
819 "Generating fill reports: instrument_id={:?}",
820 cmd.instrument_id
821 );
822
823 let account_id = self.core.account_id;
824 let start = cmd.start.map(Timestamp::from);
825 let end = cmd.end.map(Timestamp::from);
826 let mut reports = self
827 .http
828 .request_fill_reports(account_id, cmd.instrument_id, start, end)
829 .await?;
830
831 if let Some(venue_order_id) = cmd.venue_order_id {
832 reports.retain(|report| report.venue_order_id == venue_order_id);
833 }
834
835 Ok(reports)
836 }
837
838 async fn generate_position_status_reports(
839 &self,
840 cmd: &GeneratePositionStatusReports,
841 ) -> anyhow::Result<Vec<PositionStatusReport>> {
842 log::debug!(
843 "Generating position status reports: instrument_id={:?}",
844 cmd.instrument_id
845 );
846
847 let account_id = self.core.account_id;
848 self.http
849 .request_position_status_reports(account_id, cmd.instrument_id)
850 .await
851 }
852
853 async fn generate_mass_status(
854 &self,
855 lookback_mins: Option<u64>,
856 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
857 log::debug!("Generating mass status: lookback_mins={lookback_mins:?}");
858
859 let start = lookback_mins.map(|mins| Timestamp::now() - Duration::from_secs(mins * 60));
860
861 let account_id = self.core.account_id;
862 let order_reports = self
863 .http
864 .request_order_status_reports(account_id, None, start, None, true)
865 .await?;
866 let fill_reports = self
867 .http
868 .request_fill_reports(account_id, None, start, None)
869 .await?;
870 let position_reports = self
871 .http
872 .request_position_status_reports(account_id, None)
873 .await?;
874
875 let mut mass_status = ExecutionMassStatus::new(
876 self.core.client_id,
877 self.core.account_id,
878 *KRAKEN_VENUE,
879 self.clock.get_time_ns(),
880 None,
881 );
882 mass_status.add_order_reports(order_reports);
883 mass_status.add_fill_reports(fill_reports);
884 mass_status.add_position_reports(position_reports);
885
886 Ok(Some(mass_status))
887 }
888
889 fn query_account(&self, cmd: QueryAccount) -> anyhow::Result<()> {
890 log::debug!("Querying account: {cmd:?}");
891
892 let account_id = self.core.account_id;
893 let http = self.http.clone();
894 let emitter = self.emitter.clone();
895
896 self.spawn_task("query_account", async move {
897 let account_state = http.request_account_state(account_id).await?;
898 emitter.emit_account_state(
899 account_state.balances.clone(),
900 account_state.margins.clone(),
901 account_state.is_reported,
902 account_state.ts_event,
903 account_state.info,
904 );
905 Ok(())
906 });
907
908 Ok(())
909 }
910
911 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
912 log::debug!("Querying order: {cmd:?}");
913
914 let venue_order_id = cmd
915 .venue_order_id
916 .context("venue_order_id required for query_order")?;
917 let account_id = self.core.account_id;
918 let http = self.http.clone();
919 let emitter = self.emitter.clone();
920
921 self.spawn_task("query_order", async move {
922 let reports = http
923 .request_order_status_reports(account_id, None, None, None, true)
924 .await
925 .context("Failed to query order")?;
926
927 if let Some(report) = reports
928 .into_iter()
929 .find(|r| r.venue_order_id == venue_order_id)
930 {
931 emitter.send_order_status_report(report);
932 }
933 Ok(())
934 });
935
936 Ok(())
937 }
938
939 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
940 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
941 self.submit_single_order(&order, "submit_order");
942 Ok(())
943 }
944
945 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
946 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
947
948 log::debug!(
949 "Submitting order list: order_list_id={}, count={}",
950 cmd.order_list.id,
951 orders.len()
952 );
953
954 let mut order_tuples = Vec::with_capacity(orders.len());
955 let mut order_meta = Vec::with_capacity(orders.len());
956
957 for order in &orders {
958 if order.is_closed() {
959 log::warn!(
960 "Cannot submit closed order: client_order_id={}",
961 order.client_order_id()
962 );
963 continue;
964 }
965
966 if order.order_type() == OrderType::Market {
969 self.submit_single_order(order, "submit_order_list");
970 continue;
971 }
972
973 let client_order_id = order.client_order_id();
974 let kraken_cl_ord_id = truncate_cl_ord_id(&client_order_id);
975
976 if kraken_cl_ord_id != client_order_id.as_str() {
977 self.truncated_id_map
978 .insert(kraken_cl_ord_id, client_order_id);
979 }
980
981 self.register_order_identity(order);
982 self.emitter.emit_order_submitted(order);
983
984 order_tuples.push((
985 order.instrument_id(),
986 client_order_id,
987 order.order_side(),
988 order.order_type(),
989 order.quantity(),
990 order.time_in_force(),
991 order.price(),
992 order.trigger_price(),
993 order.trigger_type(),
994 order.is_reduce_only(),
995 order.is_post_only(),
996 ));
997
998 order_meta.push((order.strategy_id(), order.instrument_id(), client_order_id));
999 }
1000
1001 if order_tuples.is_empty() {
1002 return Ok(());
1003 }
1004
1005 let http = self.http.clone();
1006 let emitter = self.emitter.clone();
1007 let clock = self.clock;
1008 let dispatch_state = self.ws_dispatch_state.clone();
1009
1010 self.spawn_task("submit_order_list", async move {
1011 let results = http.send_order_batches(order_tuples).await;
1012 for (result, (strategy_id, instrument_id, client_order_id)) in
1013 results.into_iter().zip(&order_meta)
1014 {
1015 let outcome = match result {
1016 Ok(item) => command_failure_from_futures_batch_item(item),
1017 Err(e) => Err(command_failure_from_futures_batch_error(&e)),
1018 };
1019
1020 match outcome {
1021 Ok(()) => {}
1022 Err(CommandFailure::Ambiguous(reason)) => {
1023 log::warn!(
1024 "submit_order_list outcome is ambiguous for client_order_id={client_order_id}: {reason}"
1025 );
1026 }
1027 Err(
1028 CommandFailure::NotSent(reason)
1029 | CommandFailure::VenueRejected(reason),
1030 ) => {
1031 let ts_event = clock.get_time_ns();
1032 let error_msg =
1033 format!("submit_order_list batch item rejected: {reason}");
1034 dispatch_state.cleanup_terminal(client_order_id);
1035 emitter.emit_order_rejected_event(
1036 *strategy_id,
1037 *instrument_id,
1038 *client_order_id,
1039 &error_msg,
1040 ts_event,
1041 reason == "postWouldExecute",
1042 );
1043 }
1044 }
1045 }
1046 Ok(())
1047 });
1048
1049 Ok(())
1050 }
1051
1052 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1053 self.modify_single_order(&cmd);
1054 Ok(())
1055 }
1056
1057 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1058 self.cancel_single_order(&cmd);
1059 Ok(())
1060 }
1061
1062 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
1063 let instrument_id = cmd.instrument_id;
1064
1065 if cmd.order_side.is_none() {
1066 log::debug!("Canceling all orders: instrument_id={instrument_id} (bulk)");
1067
1068 let http = self.http.clone();
1069 let symbol = instrument_id.symbol.to_string();
1070
1071 self.spawn_task("cancel_all_orders", async move {
1072 match http.inner.cancel_all_orders(Some(symbol)).await {
1073 Ok(response) => {
1074 if response.result != KrakenApiResult::Success
1075 && response.cancel_status.cancelled_orders.is_empty()
1076 {
1077 log::warn!(
1078 "Cancel-all failed without per-order results, awaiting reconciliation: status={}",
1079 response.cancel_status.status
1080 );
1081 }
1082 }
1083 Err(e) => match command_failure_from_cancel_error(e) {
1084 CommandFailure::NotSent(reason) => {
1085 log::warn!("Cancel-all failed local validation: {reason}");
1086 }
1087 CommandFailure::Ambiguous(reason)
1088 | CommandFailure::VenueRejected(reason) => {
1089 log::warn!(
1090 "Cancel-all ambiguous failure, awaiting reconciliation: {reason}"
1091 );
1092 }
1093 },
1094 }
1095 Ok(())
1096 });
1097
1098 return Ok(());
1099 }
1100
1101 log::debug!(
1102 "Canceling all orders: instrument_id={instrument_id}, side={:?}",
1103 cmd.order_side
1104 );
1105
1106 let orders_to_cancel: Vec<_> = {
1107 let cache = self.core.cache();
1108 let open_orders = cache.orders_open(None, Some(&instrument_id), None, None, None);
1109
1110 open_orders
1111 .into_iter()
1112 .filter(|order| Some(order.order_side()) == cmd.order_side)
1113 .filter_map(|order| {
1114 Some((
1115 order.venue_order_id()?,
1116 order.client_order_id(),
1117 order.instrument_id(),
1118 order.strategy_id(),
1119 ))
1120 })
1121 .collect()
1122 };
1123
1124 let account_id = self.core.account_id;
1125
1126 for (venue_order_id, client_order_id, order_instrument_id, strategy_id) in orders_to_cancel
1127 {
1128 let http = self.http.clone();
1129 let emitter = self.emitter.clone();
1130 let clock = self.clock;
1131
1132 self.spawn_task("cancel_order_by_side", async move {
1133 if let Err(failure) = cancel_order_for_futures(
1134 &http,
1135 account_id,
1136 order_instrument_id,
1137 Some(client_order_id),
1138 Some(venue_order_id),
1139 )
1140 .await
1141 {
1142 handle_cancel_failure(
1143 &emitter,
1144 clock,
1145 strategy_id,
1146 order_instrument_id,
1147 client_order_id,
1148 Some(venue_order_id),
1149 failure,
1150 );
1151 }
1152 Ok(())
1153 });
1154 }
1155
1156 Ok(())
1157 }
1158
1159 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1160 log::debug!(
1161 "Batch canceling orders: instrument_id={}, count={}",
1162 cmd.instrument_id,
1163 cmd.cancels.len()
1164 );
1165
1166 let http = self.http.clone();
1167 let emitter = self.emitter.clone();
1168 let clock = self.clock;
1169 let cancels = cmd.cancels;
1170
1171 self.spawn_task("batch_cancel_orders", async move {
1172 batch_cancel_orders_for_futures(&http, &emitter, clock, &cancels).await;
1173 Ok(())
1174 });
1175
1176 Ok(())
1177 }
1178}
1179
1180#[derive(Debug, Clone)]
1181struct CancelRequestContext {
1182 strategy_id: StrategyId,
1183 instrument_id: InstrumentId,
1184 client_order_id: ClientOrderId,
1185 truncated_client_order_id: String,
1186 venue_order_id: Option<VenueOrderId>,
1187}
1188
1189async fn cancel_order_for_futures(
1190 http: &KrakenFuturesHttpClient,
1191 _account_id: AccountId,
1192 instrument_id: InstrumentId,
1193 client_order_id: Option<ClientOrderId>,
1194 venue_order_id: Option<VenueOrderId>,
1195) -> Result<(), CommandFailure> {
1196 http.get_cached_instrument(&instrument_id.symbol.inner())
1197 .ok_or_else(|| {
1198 CommandFailure::not_sent(InstrumentLookupError::not_found(instrument_id).to_string())
1199 })?;
1200
1201 let order_id = venue_order_id.as_ref().map(ToString::to_string);
1202 let cli_ord_id = client_order_id.as_ref().map(truncate_cl_ord_id);
1203
1204 if order_id.is_none() && cli_ord_id.is_none() {
1205 return Err(CommandFailure::not_sent(
1206 "Either client_order_id or venue_order_id must be provided",
1207 ));
1208 }
1209
1210 let response = http
1211 .inner
1212 .cancel_order(order_id, cli_ord_id)
1213 .await
1214 .map_err(command_failure_from_cancel_error)?;
1215
1216 if response.result != KrakenApiResult::Success
1217 || response.cancel_status.status != KrakenSendStatus::Cancelled
1218 {
1219 return Err(CommandFailure::venue_rejected(format!(
1220 "cancel-order rejected: status={}",
1221 response.cancel_status.status
1222 )));
1223 }
1224
1225 Ok(())
1226}
1227
1228async fn batch_cancel_orders_for_futures(
1229 http: &KrakenFuturesHttpClient,
1230 emitter: &ExecutionEventEmitter,
1231 clock: &'static AtomicTime,
1232 cancels: &[CancelOrder],
1233) {
1234 let mut contexts = Vec::new();
1235 let mut items = Vec::new();
1236
1237 for cancel in cancels {
1238 match batch_cancel_item_for_futures(http, cancel) {
1239 Ok((context, item)) => {
1240 contexts.push(context);
1241 items.push(item);
1242 }
1243 Err(CommandFailure::NotSent(reason)) => {
1244 log::warn!(
1245 "Batch cancel command failed local validation for {}: {reason}",
1246 cancel.client_order_id
1247 );
1248 }
1249 Err(CommandFailure::Ambiguous(reason) | CommandFailure::VenueRejected(reason)) => {
1250 log::warn!(
1251 "Batch cancel command ambiguous failure for {}, awaiting reconciliation: {reason}",
1252 cancel.client_order_id
1253 );
1254 }
1255 }
1256 }
1257
1258 for (item_chunk, context_chunk) in items
1259 .chunks(FUTURES_BATCH_CANCEL_LIMIT)
1260 .zip(contexts.chunks(FUTURES_BATCH_CANCEL_LIMIT))
1261 {
1262 let response = match http
1263 .inner
1264 .cancel_order_items_batch(item_chunk.to_vec())
1265 .await
1266 {
1267 Ok(response) => response,
1268 Err(e) => {
1269 match command_failure_from_cancel_error(e) {
1270 CommandFailure::NotSent(reason) => {
1271 log::warn!("Batch cancel failed local validation: {reason}");
1272 }
1273 CommandFailure::Ambiguous(reason) | CommandFailure::VenueRejected(reason) => {
1274 log::warn!(
1275 "Batch cancel failed without per-order results, awaiting reconciliation: {reason}"
1276 );
1277 }
1278 }
1279 continue;
1280 }
1281 };
1282
1283 if response.batch_status.is_empty() {
1284 if response.result != KrakenApiResult::Success {
1285 let reason = response.error.as_deref().unwrap_or("Unknown error");
1286 log::warn!(
1287 "Batch cancel failed without per-order results, awaiting reconciliation: {reason}"
1288 );
1289 }
1290 continue;
1291 }
1292
1293 if response.batch_status.len() != context_chunk.len() {
1294 log::warn!(
1295 "Batch cancel returned {} per-order result(s) for {} request(s); unmatched results await reconciliation",
1296 response.batch_status.len(),
1297 context_chunk.len()
1298 );
1299 }
1300
1301 for (index, status) in response.batch_status.iter().enumerate() {
1302 let Some(cancel_status) = batch_cancel_status(status) else {
1303 log::warn!("Batch cancel result without status at index {index}");
1304 continue;
1305 };
1306
1307 if cancel_status == KrakenSendStatus::Cancelled {
1308 continue;
1309 }
1310
1311 let Some(context) = batch_cancel_context(status, context_chunk, index) else {
1312 log::warn!(
1313 "Batch cancel rejected item without matching request context at index {index}: status={cancel_status}"
1314 );
1315 continue;
1316 };
1317
1318 emitter.emit_order_cancel_rejected_event(
1319 context.strategy_id,
1320 context.instrument_id,
1321 context.client_order_id,
1322 context.venue_order_id,
1323 &format!("batch-cancel rejected: status={cancel_status}"),
1324 clock.get_time_ns(),
1325 );
1326 }
1327 }
1328}
1329
1330fn batch_cancel_item_for_futures(
1331 http: &KrakenFuturesHttpClient,
1332 cancel: &CancelOrder,
1333) -> Result<(CancelRequestContext, KrakenFuturesBatchCancelItem), CommandFailure> {
1334 http.get_cached_instrument(&cancel.instrument_id.symbol.inner())
1335 .ok_or_else(|| {
1336 CommandFailure::not_sent(
1337 InstrumentLookupError::not_found(cancel.instrument_id).to_string(),
1338 )
1339 })?;
1340
1341 let truncated_client_order_id = truncate_cl_ord_id(&cancel.client_order_id);
1342 let item = if let Some(venue_order_id) = cancel.venue_order_id {
1343 KrakenFuturesBatchCancelItem::from_order_id(venue_order_id.to_string())
1344 } else {
1345 KrakenFuturesBatchCancelItem::from_client_order_id(truncated_client_order_id.clone())
1346 };
1347
1348 Ok((
1349 CancelRequestContext {
1350 strategy_id: cancel.strategy_id,
1351 instrument_id: cancel.instrument_id,
1352 client_order_id: cancel.client_order_id,
1353 truncated_client_order_id,
1354 venue_order_id: cancel.venue_order_id,
1355 },
1356 item,
1357 ))
1358}
1359
1360fn batch_cancel_status(status: &FuturesBatchCancelStatus) -> Option<KrakenSendStatus> {
1361 status
1362 .cancel_status
1363 .as_ref()
1364 .map(|cancel_status| cancel_status.status)
1365 .or(status.status)
1366}
1367
1368fn batch_cancel_context<'a>(
1369 status: &FuturesBatchCancelStatus,
1370 contexts: &'a [CancelRequestContext],
1371 index: usize,
1372) -> Option<&'a CancelRequestContext> {
1373 if let Some(order_id) = status.order_id.as_deref()
1374 && let Some(context) = contexts.iter().find(|context| {
1375 context
1376 .venue_order_id
1377 .is_some_and(|venue_order_id| venue_order_id.as_str() == order_id)
1378 })
1379 {
1380 return Some(context);
1381 }
1382
1383 if let Some(cli_ord_id) = status.cli_ord_id.as_deref()
1384 && let Some(context) = contexts
1385 .iter()
1386 .find(|context| context.truncated_client_order_id == cli_ord_id)
1387 {
1388 return Some(context);
1389 }
1390
1391 if index < contexts.len() && status.order_id.is_none() && status.cli_ord_id.is_none() {
1392 return contexts.get(index);
1393 }
1394
1395 None
1396}
1397
1398fn handle_cancel_failure(
1399 emitter: &ExecutionEventEmitter,
1400 clock: &'static AtomicTime,
1401 strategy_id: StrategyId,
1402 instrument_id: InstrumentId,
1403 client_order_id: ClientOrderId,
1404 venue_order_id: Option<VenueOrderId>,
1405 failure: CommandFailure,
1406) {
1407 match failure {
1408 CommandFailure::VenueRejected(reason) => {
1409 emitter.emit_order_cancel_rejected_event(
1410 strategy_id,
1411 instrument_id,
1412 client_order_id,
1413 venue_order_id,
1414 &reason,
1415 clock.get_time_ns(),
1416 );
1417 }
1418 CommandFailure::NotSent(reason) => {
1419 log::warn!("Cancel command failed local validation for {client_order_id}: {reason}");
1420 }
1421 CommandFailure::Ambiguous(reason) => {
1422 log::warn!(
1423 "Ambiguous cancel failure for {client_order_id}, awaiting reconciliation: {reason}"
1424 );
1425 }
1426 }
1427}
1428
1429impl KrakenFuturesExecutionClient {
1430 fn get_cached_order_for_status_command(
1431 &self,
1432 cmd: &GenerateOrderStatusReport,
1433 ) -> Option<OrderAny> {
1434 let cache = self.core.cache();
1435
1436 if let Some(client_order_id) = cmd.client_order_id {
1437 return cache.order(&client_order_id).map(|o| o.clone());
1438 }
1439
1440 let venue_order_id = cmd.venue_order_id?;
1441 let client_order_id = *cache.client_order_id(&venue_order_id)?;
1442 cache.order(&client_order_id).map(|o| o.clone())
1443 }
1444}
1445
1446fn synthesize_filled_order_status_report(
1447 cmd: &GenerateOrderStatusReport,
1448 order: &OrderAny,
1449 fills: &[FillReport],
1450) -> Option<OrderStatusReport> {
1451 let venue_order_id = cmd.venue_order_id.or(order.venue_order_id());
1452 let truncated_client_order_id = truncate_cl_ord_id(&order.client_order_id());
1453
1454 let mut matched: Vec<&FillReport> = if let Some(venue_order_id) = venue_order_id {
1455 fills
1456 .iter()
1457 .filter(|fill| fill.venue_order_id == venue_order_id)
1458 .collect()
1459 } else {
1460 Vec::new()
1461 };
1462
1463 if matched.is_empty() {
1464 matched = fills
1465 .iter()
1466 .filter(|fill| {
1467 fill.client_order_id == Some(order.client_order_id())
1468 || fill
1469 .client_order_id
1470 .as_ref()
1471 .is_some_and(|fill_client_order_id| {
1472 fill_client_order_id.as_str() == truncated_client_order_id
1473 })
1474 })
1475 .collect();
1476 }
1477
1478 if matched.is_empty() {
1479 return None;
1480 }
1481
1482 matched.sort_by_key(|fill| fill.ts_event);
1483 let first_fill = *matched.first()?;
1484 let last_fill = *matched.last()?;
1485
1486 let total_filled = matched
1487 .iter()
1488 .fold(Decimal::ZERO, |acc, fill| acc + fill.last_qty.as_decimal());
1489 if total_filled < order.quantity().as_decimal() {
1490 return None;
1491 }
1492
1493 let total_notional = matched.iter().fold(Decimal::ZERO, |acc, fill| {
1494 acc + fill.last_qty.as_decimal() * fill.last_px.as_decimal()
1495 });
1496 let avg_px = if total_filled.is_zero() {
1497 None
1498 } else {
1499 Some(total_notional / total_filled)
1500 };
1501 let venue_order_id = venue_order_id.unwrap_or(first_fill.venue_order_id);
1502
1503 let mut report = OrderStatusReport::new(
1504 first_fill.account_id,
1505 order.instrument_id(),
1506 Some(order.client_order_id()),
1507 venue_order_id,
1508 order.order_side().into(),
1509 order.order_type(),
1510 order.time_in_force(),
1511 OrderStatus::Filled,
1512 order.quantity(),
1513 order.quantity(),
1514 first_fill.ts_event,
1515 last_fill.ts_event,
1516 last_fill.ts_init,
1517 None,
1518 );
1519 report.order_list_id = order.order_list_id();
1520 report.venue_position_id = matched.iter().rev().find_map(|fill| fill.venue_position_id);
1521 report.linked_order_ids = order
1522 .linked_order_ids()
1523 .map(|linked_order_ids| linked_order_ids.to_vec());
1524 report.parent_order_id = order.parent_order_id();
1525 report.expire_time = order.expire_time();
1526 report.price = order.price();
1527 report.trigger_price = order.trigger_price();
1528 report.trigger_type = order.trigger_type();
1529 report.avg_px = avg_px;
1530 report.display_qty = order.display_qty();
1531 report.post_only = order.is_post_only();
1532 report.reduce_only = order.is_reduce_only();
1533 Some(report)
1534}
1535
1536#[cfg(test)]
1537mod tests {
1538 use nautilus_core::{UUID4, UnixNanos};
1539 use nautilus_model::{
1540 enums::{LiquiditySide, OrderSide, OrderType, TimeInForce},
1541 identifiers::{
1542 AccountId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, VenueOrderId,
1543 },
1544 orders::OrderTestBuilder,
1545 reports::FillReport,
1546 types::{Currency, Money, Price, Quantity},
1547 };
1548 use rstest::rstest;
1549
1550 use super::*;
1551
1552 const TEST_INSTRUMENT_ID: &str = "PF_XBTUSD.KRAKEN";
1553
1554 #[tokio::test]
1555 async fn test_cancel_order_for_futures_missing_cached_instrument_returns_canonical_error() {
1556 let http = KrakenFuturesHttpClient::default();
1557 let instrument_id = InstrumentId::from(TEST_INSTRUMENT_ID);
1558 let client_order_id = ClientOrderId::from("C-001");
1559 let venue_order_id = VenueOrderId::from("V-001");
1560
1561 let result = cancel_order_for_futures(
1562 &http,
1563 AccountId::from("KRAKEN-001"),
1564 instrument_id,
1565 Some(client_order_id),
1566 Some(venue_order_id),
1567 )
1568 .await;
1569
1570 match result {
1571 Err(CommandFailure::NotSent(reason)) => {
1572 assert_eq!(
1573 reason,
1574 InstrumentLookupError::not_found(instrument_id).to_string()
1575 );
1576 }
1577 _ => panic!("Expected local validation failure"),
1578 }
1579 }
1580
1581 #[rstest]
1582 fn test_batch_cancel_item_for_futures_missing_cached_instrument_returns_canonical_error() {
1583 let http = KrakenFuturesHttpClient::default();
1584 let instrument_id = InstrumentId::from(TEST_INSTRUMENT_ID);
1585 let cancel = CancelOrder::new(
1586 TraderId::from("TESTER-001"),
1587 None,
1588 StrategyId::from("S-001"),
1589 instrument_id,
1590 ClientOrderId::from("C-001"),
1591 Some(VenueOrderId::from("V-001")),
1592 UUID4::new(),
1593 UnixNanos::default(),
1594 None,
1595 None,
1596 );
1597
1598 let result = batch_cancel_item_for_futures(&http, &cancel);
1599
1600 match result {
1601 Err(CommandFailure::NotSent(reason)) => {
1602 assert_eq!(
1603 reason,
1604 InstrumentLookupError::not_found(instrument_id).to_string()
1605 );
1606 }
1607 _ => panic!("Expected local validation failure"),
1608 }
1609 }
1610
1611 fn make_fill(
1612 venue_order_id: &str,
1613 client_order_id: Option<&str>,
1614 quantity: &str,
1615 price: &str,
1616 ts_event: u64,
1617 ) -> FillReport {
1618 FillReport::new(
1619 AccountId::from("KRAKEN-001"),
1620 InstrumentId::from(TEST_INSTRUMENT_ID),
1621 VenueOrderId::from(venue_order_id),
1622 TradeId::from(format!("T-{ts_event}").as_str()),
1623 OrderSide::Buy,
1624 Quantity::from(quantity),
1625 Price::from(price),
1626 Money::from_decimal(Decimal::ZERO, Currency::USD()).unwrap(),
1627 LiquiditySide::Taker,
1628 client_order_id.map(ClientOrderId::from),
1629 None,
1630 UnixNanos::from(ts_event),
1631 UnixNanos::from(ts_event),
1632 None,
1633 )
1634 }
1635
1636 fn make_cmd(
1637 client_order_id: Option<&str>,
1638 venue_order_id: Option<&str>,
1639 ) -> GenerateOrderStatusReport {
1640 GenerateOrderStatusReport::new(
1641 UUID4::new(),
1642 UnixNanos::default(),
1643 Some(InstrumentId::from(TEST_INSTRUMENT_ID)),
1644 client_order_id.map(ClientOrderId::from),
1645 venue_order_id.map(VenueOrderId::from),
1646 None,
1647 None,
1648 )
1649 }
1650
1651 fn make_order(client_order_id: &str) -> OrderAny {
1652 OrderTestBuilder::new(OrderType::Market)
1653 .instrument_id(InstrumentId::from(TEST_INSTRUMENT_ID))
1654 .client_order_id(ClientOrderId::from(client_order_id))
1655 .side(OrderSide::Buy)
1656 .quantity(Quantity::from("100"))
1657 .time_in_force(TimeInForce::Ioc)
1658 .build()
1659 }
1660
1661 #[rstest]
1662 fn test_synthesize_filled_order_status_report_matches_full_fill_by_venue_order_id() {
1663 let order = make_order("O-123456");
1664 let cmd = make_cmd(Some("O-123456"), Some("KRAKEN-789"));
1665 let fills = vec![
1666 make_fill("KRAKEN-789", Some("O-123456"), "40", "50000.0", 1),
1667 make_fill("KRAKEN-789", Some("O-123456"), "60", "50010.0", 2),
1668 make_fill("KRAKEN-OTHER", Some("O-123456"), "999", "1.0", 3),
1669 ];
1670
1671 let report = synthesize_filled_order_status_report(&cmd, &order, &fills)
1672 .expect("expected a filled report");
1673
1674 assert_eq!(report.venue_order_id, VenueOrderId::from("KRAKEN-789"));
1675 assert_eq!(
1676 report.client_order_id,
1677 Some(ClientOrderId::from("O-123456"))
1678 );
1679 assert_eq!(report.order_status, OrderStatus::Filled);
1680 assert_eq!(report.order_type, OrderType::Market);
1681 assert_eq!(report.time_in_force, TimeInForce::Ioc);
1682 assert_eq!(report.quantity, Quantity::from("100"));
1683 assert_eq!(report.filled_qty, Quantity::from("100"));
1684 assert_eq!(
1685 report.avg_px,
1686 Some(Decimal::from_str_exact("50006.0").unwrap())
1687 );
1688 }
1689
1690 #[rstest]
1691 fn test_synthesize_filled_order_status_report_requires_full_fill_size() {
1692 let order = make_order("O-123457");
1693 let cmd = make_cmd(Some("O-123457"), Some("KRAKEN-790"));
1694 let fills = vec![make_fill(
1695 "KRAKEN-790",
1696 Some("O-123457"),
1697 "40",
1698 "50000.0",
1699 1,
1700 )];
1701
1702 assert!(synthesize_filled_order_status_report(&cmd, &order, &fills).is_none());
1703 }
1704
1705 #[rstest]
1706 fn test_synthesize_filled_order_status_report_matches_truncated_client_order_id() {
1707 let long_client_order_id = "O202602270023210040011";
1708 let order = make_order(long_client_order_id);
1709 let cmd = make_cmd(Some(long_client_order_id), None);
1710 let fills = vec![make_fill(
1711 "KRAKEN-791",
1712 Some(truncate_cl_ord_id(&ClientOrderId::from(long_client_order_id)).as_str()),
1713 "100",
1714 "50000.0",
1715 1,
1716 )];
1717
1718 let report = synthesize_filled_order_status_report(&cmd, &order, &fills)
1719 .expect("expected a filled report");
1720
1721 assert_eq!(
1722 report.client_order_id,
1723 Some(ClientOrderId::from(long_client_order_id))
1724 );
1725 assert_eq!(report.venue_order_id, VenueOrderId::from("KRAKEN-791"));
1726 assert_eq!(report.order_status, OrderStatus::Filled);
1727 }
1728}