Skip to main content

nautilus_kraken/execution/
futures.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Kraken Futures execution client implementation.
17
18use 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
87/// Maximum order IDs per `/orders/status` request for Kraken Futures API.
88const FUTURES_ORDERS_STATUS_LIMIT: usize = 50;
89
90/// Kraken Futures execution client.
91///
92/// Provides order management, account operations, and position management
93/// for Kraken Futures markets.
94#[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    /// Creates a new [`KrakenFuturesExecutionClient`].
116    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    /// Returns a reference to the clock.
202    #[must_use]
203    pub fn clock(&self) -> &'static AtomicTime {
204        self.clock
205    }
206
207    /// Returns a reference to the event emitter.
208    #[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        // Match by venue_order_id or client_order_id (comparing truncated form
771        // since Kraken stores the truncated cl_ord_id for long IDs)
772        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        // Held orders never appear on /openorders; query the 5-second window
791        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        // Window filled reports have no avg_px; price them from fills below
823        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            // Unpriced filled reports would close at the order price
848            (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        // Snapshot recon would infer uncovered fills at the order price
955        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            // Kraken batch endpoint only supports limit and stop orders,
1069            // submit market orders individually
1070            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    /// Returns `/orders/status` reports for cached-open orders the given
1533    /// venue reports do not cover.
1534    ///
1535    /// A Maker Protection hold never reaches the book, so it is invisible to
1536    /// an open-orders snapshot while the venue still knows the order within
1537    /// its 5-second `/orders/status` window. Orders without a venue order ID
1538    /// yet (their submit acknowledgement is still inside the hold window) are
1539    /// queried by their truncated client order ID, the only venue handle that
1540    /// exists during the window.
1541    ///
1542    /// # Errors
1543    ///
1544    /// Returns an error if the underlying request fails or the venue rejects
1545    /// it, so callers defer rather than treat the orders as missing.
1546    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                // Spot and Futures share the KRAKEN venue and one cache
1565                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}