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