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