Skip to main content

nautilus_kraken/websocket/dispatch/
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//! WebSocket execution dispatch for the Kraken Futures API.
17//!
18//! Routes `OpenOrdersDelta`, `OpenOrdersCancel`, and `FillsDelta` messages to
19//! typed order events (for tracked orders) or status / fill reports (for
20//! external orders) under the two-tier dispatch contract.
21
22use std::sync::Arc;
23
24use ahash::AHashMap;
25use nautilus_core::{AtomicMap, UUID4, UnixNanos};
26use nautilus_live::ExecutionEventEmitter;
27use nautilus_model::{
28    enums::{OrderStatus, OrderType, TimeInForce},
29    events::{OrderCanceled, OrderEventAny, OrderUpdated},
30    identifiers::{AccountId, ClientOrderId, InstrumentId, VenueOrderId},
31    instruments::{Instrument, InstrumentAny},
32    reports::OrderStatusReport,
33    types::{Price, Quantity},
34};
35use ustr::Ustr;
36
37use super::{
38    DeltaSnapshot, OrderIdentity, PendingRemoval, WsDispatchState, ensure_accepted_emitted,
39    fill_report_to_order_filled, resolve_client_order_id,
40};
41use crate::{
42    common::lookup_instrument_in_snapshot,
43    websocket::futures::{
44        messages::{
45            KrakenFuturesFill, KrakenFuturesFillsDelta, KrakenFuturesOpenOrdersCancel,
46            KrakenFuturesOpenOrdersDelta,
47        },
48        parse::{parse_futures_ws_fill_report, parse_futures_ws_order_status_report},
49    },
50};
51
52/// Dispatches a Kraken Futures `OpenOrdersDelta` message.
53///
54/// Fill-driven cancel deltas (`is_cancel=true` with reason `full_fill`) are
55/// skipped - the corresponding `FillsDelta` carries the real fill, so emitting
56/// a synthetic Canceled here would race with the genuine `OrderFilled`.
57/// Part-fill removals (`is_cancel=true` with reason `partial_fill`) discard the
58/// order's remainder and are terminal: they close the order only once the
59/// fills feed has accounted the removal's cumulative filled, so a fill still
60/// in flight is never orphaned.
61#[expect(clippy::too_many_arguments)]
62pub fn open_orders_delta(
63    delta: &KrakenFuturesOpenOrdersDelta,
64    state: &WsDispatchState,
65    emitter: &ExecutionEventEmitter,
66    instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
67    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
68    order_instrument_map: &Arc<AtomicMap<String, InstrumentId>>,
69    venue_client_map: &Arc<AtomicMap<String, ClientOrderId>>,
70    venue_order_qty: &Arc<AtomicMap<String, Quantity>>,
71    account_id: AccountId,
72    ts_init: UnixNanos,
73) {
74    if delta.is_fill_driven_cancel() && !delta.is_partial_fill_removal() {
75        log::debug!(
76            "Skipping fill-driven open_orders delta: order_id={}, reason={:?}",
77            delta.order.order_id,
78            delta.reason,
79        );
80        return;
81    }
82
83    let product_id = delta.order.instrument.as_str();
84    let instruments = instruments.load();
85    let Some(instrument) = lookup_instrument_in_snapshot(&instruments, product_id) else {
86        log::warn!("No instrument for product_id: {product_id}");
87        return;
88    };
89
90    // Cache instrument and qty by venue order id so cancel-only messages
91    // (which arrive without the order body) can be reconstructed for the
92    // external fallback path.
93    order_instrument_map.insert(delta.order.order_id.clone(), instrument.id());
94    let Ok(qty) = Quantity::from_decimal_dp(delta.order.qty, instrument.size_precision()) else {
95        log::error!("Failed to parse order quantity: {}", delta.order.qty);
96        return;
97    };
98    venue_order_qty.insert(delta.order.order_id.clone(), qty);
99
100    let resolved_id = delta
101        .order
102        .cli_ord_id
103        .as_ref()
104        .map(|id| resolve_client_order_id(id, truncated_id_map));
105
106    // Stale-report suppression: an order that already reached the filled
107    // terminal state should not produce more events even if a late delta
108    // arrives. `filled_orders` persists past `cleanup_terminal` precisely
109    // for this check.
110    if let Some(cid) = resolved_id
111        && state.filled_orders.contains(&cid)
112    {
113        log::debug!(
114            "Skipping stale open_orders delta for filled order: cid={cid}, order_id={}",
115            delta.order.order_id,
116        );
117        return;
118    }
119
120    if delta.is_partial_fill_removal() {
121        // Untracked removals must not emit a terminal report; that orphans the fill
122        if let Some(client_order_id) = resolved_id {
123            venue_client_map.insert(delta.order.order_id.clone(), client_order_id);
124
125            if let Some(identity) = state.lookup_identity(&client_order_id) {
126                partial_removal_tracked(
127                    delta,
128                    client_order_id,
129                    &identity,
130                    instrument,
131                    state,
132                    emitter,
133                    account_id,
134                    ts_init,
135                );
136                return;
137            }
138        }
139
140        log::debug!(
141            "Skipping untracked partial-fill removal: order_id={}",
142            delta.order.order_id,
143        );
144        return;
145    }
146
147    if let Some(client_order_id) = resolved_id {
148        venue_client_map.insert(delta.order.order_id.clone(), client_order_id);
149
150        if let Some(identity) = state.lookup_identity(&client_order_id) {
151            delta_tracked(
152                delta,
153                client_order_id,
154                &identity,
155                instrument,
156                state,
157                emitter,
158                account_id,
159                ts_init,
160            );
161            return;
162        }
163    }
164
165    // External / untracked: fall back to a status report.
166    match parse_futures_ws_order_status_report(
167        &delta.order,
168        delta.is_cancel,
169        delta.reason.as_deref(),
170        instrument,
171        account_id,
172        ts_init,
173    ) {
174        Ok(mut report) => {
175            if let Some(cid) = resolved_id {
176                report = report.with_client_order_id(cid);
177            }
178            emitter.send_order_status_report(report);
179        }
180        Err(e) => log::error!("Failed to parse futures order status report: {e}"),
181    }
182}
183
184#[expect(clippy::too_many_arguments)]
185fn delta_tracked(
186    delta: &KrakenFuturesOpenOrdersDelta,
187    client_order_id: ClientOrderId,
188    identity: &OrderIdentity,
189    instrument: &InstrumentAny,
190    state: &WsDispatchState,
191    emitter: &ExecutionEventEmitter,
192    account_id: AccountId,
193    ts_init: UnixNanos,
194) {
195    let venue_order_id = VenueOrderId::new(&delta.order.order_id);
196    let ts_event = millis_to_nanos(delta.order.last_update_time);
197    let Ok(new_filled) = Quantity::from_decimal_dp(delta.order.filled, instrument.size_precision())
198    else {
199        log::error!("Failed to parse filled quantity: {}", delta.order.filled);
200        return;
201    };
202
203    if delta.is_cancel {
204        ensure_accepted_emitted(
205            client_order_id,
206            venue_order_id,
207            account_id,
208            identity,
209            state,
210            emitter,
211            ts_event,
212            ts_init,
213        );
214        let canceled = OrderCanceled::new(
215            emitter.trader_id(),
216            identity.strategy_id,
217            identity.instrument_id,
218            client_order_id,
219            UUID4::new(),
220            ts_event,
221            ts_init,
222            false,
223            Some(venue_order_id),
224            Some(account_id),
225            delta.reason.as_deref().map(Ustr::from),
226        );
227        emitter.send_order_event(OrderEventAny::Canceled(canceled));
228        state.cleanup_terminal(&client_order_id);
229        return;
230    }
231
232    let already_accepted = state.emitted_accepted.contains(&client_order_id);
233    ensure_accepted_emitted(
234        client_order_id,
235        venue_order_id,
236        account_id,
237        identity,
238        state,
239        emitter,
240        ts_event,
241        ts_init,
242    );
243
244    let Ok(qty) = Quantity::from_decimal_dp(delta.order.qty, instrument.size_precision()) else {
245        log::error!("Failed to parse order quantity: {}", delta.order.qty);
246        return;
247    };
248    let snapshot = DeltaSnapshot::new(
249        qty,
250        new_filled,
251        delta.order.limit_price,
252        delta.order.stop_price,
253    );
254
255    if !already_accepted {
256        // First delta seen for this order: the placement Accepted is enough.
257        state.record_delta_snapshot(client_order_id, snapshot);
258        return;
259    }
260
261    // Follow-up delta. The two emission-relevant signals are independent:
262    //   * filled increased       -> partial-fill notification (FillsDelta has it,
263    //                               nothing to emit from here)
264    //   * non-fill field changed -> modify acknowledgement (emit OrderUpdated)
265    // Both can be true simultaneously when a user amends a partially filled
266    // order, so check the modify branch regardless of fill movement.
267    let previous = state.previous_delta_snapshot(&client_order_id);
268    state.record_delta_snapshot(client_order_id, snapshot);
269
270    let non_fill_changed = previous.is_some_and(|prev| !snapshot.non_fill_fields_match(&prev));
271    if !non_fill_changed {
272        return;
273    }
274
275    // Modify ack: refresh tracked quantity (size may have changed) and emit
276    // OrderUpdated so the engine clears PendingUpdate.
277    state.update_identity_quantity(&client_order_id, qty);
278    let price = match delta
279        .order
280        .limit_price
281        .map(|p| Price::from_decimal_dp(p, instrument.price_precision()))
282        .transpose()
283    {
284        Ok(price) => price,
285        Err(e) => {
286            log::error!("Failed to parse limit price: {e}");
287            return;
288        }
289    };
290    let trigger_price = match delta
291        .order
292        .stop_price
293        .map(|p| Price::from_decimal_dp(p, instrument.price_precision()))
294        .transpose()
295    {
296        Ok(price) => price,
297        Err(e) => {
298            log::error!("Failed to parse stop price: {e}");
299            return;
300        }
301    };
302
303    let updated = OrderUpdated::new(
304        emitter.trader_id(),
305        identity.strategy_id,
306        identity.instrument_id,
307        client_order_id,
308        qty,
309        UUID4::new(),
310        ts_event,
311        ts_init,
312        false,
313        Some(venue_order_id),
314        Some(account_id),
315        price,
316        trigger_price,
317        None,
318        false,
319    );
320    emitter.send_order_event(OrderEventAny::Updated(updated));
321}
322
323/// Converges a tracked terminal part-fill removal whose remainder the venue
324/// discarded (a converted Maker Protection hold or an IOC-style order).
325///
326/// The removal delta carries the venue's cumulative filled at the removal. If
327/// the fills feed has already accounted it the order closes immediately;
328/// otherwise the removal is parked until the cumulative fills reach it, so
329/// the closing cancel can never overtake the fill it belongs with.
330#[expect(clippy::too_many_arguments)]
331fn partial_removal_tracked(
332    delta: &KrakenFuturesOpenOrdersDelta,
333    client_order_id: ClientOrderId,
334    identity: &OrderIdentity,
335    instrument: &InstrumentAny,
336    state: &WsDispatchState,
337    emitter: &ExecutionEventEmitter,
338    account_id: AccountId,
339    ts_init: UnixNanos,
340) {
341    let venue_order_id = VenueOrderId::new(&delta.order.order_id);
342    let ts_event = millis_to_nanos(delta.order.last_update_time);
343
344    let Ok(venue_filled) =
345        Quantity::from_decimal_dp(delta.order.filled, instrument.size_precision())
346    else {
347        log::error!("Failed to parse filled quantity: {}", delta.order.filled);
348        return;
349    };
350
351    let recorded_filled = state
352        .previous_filled_qty(&client_order_id)
353        .unwrap_or_else(|| Quantity::zero(instrument.size_precision()));
354
355    if recorded_filled < venue_filled {
356        log::debug!(
357            "Deferring partial-fill removal for {client_order_id}: filled={recorded_filled}, \
358             venue_filled={venue_filled}",
359        );
360        state.insert_pending_removal(
361            client_order_id,
362            PendingRemoval {
363                venue_filled,
364                reason: delta.reason.clone(),
365                ts_event,
366            },
367        );
368
369        return;
370    }
371
372    ensure_accepted_emitted(
373        client_order_id,
374        venue_order_id,
375        account_id,
376        identity,
377        state,
378        emitter,
379        ts_event,
380        ts_init,
381    );
382
383    let canceled = OrderCanceled::new(
384        emitter.trader_id(),
385        identity.strategy_id,
386        identity.instrument_id,
387        client_order_id,
388        UUID4::new(),
389        ts_event,
390        ts_init,
391        false,
392        Some(venue_order_id),
393        Some(account_id),
394        delta.reason.as_deref().map(Ustr::from),
395    );
396    emitter.send_order_event(OrderEventAny::Canceled(canceled));
397    state.cleanup_terminal(&client_order_id);
398}
399
400/// Dispatches a Kraken Futures `OpenOrdersCancel` (cancel-only) message.
401#[expect(clippy::too_many_arguments)]
402pub fn open_orders_cancel(
403    cancel: &KrakenFuturesOpenOrdersCancel,
404    state: &WsDispatchState,
405    emitter: &ExecutionEventEmitter,
406    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
407    order_instrument_map: &Arc<AtomicMap<String, InstrumentId>>,
408    venue_client_map: &Arc<AtomicMap<String, ClientOrderId>>,
409    venue_order_qty: &Arc<AtomicMap<String, Quantity>>,
410    account_id: AccountId,
411    ts_init: UnixNanos,
412) {
413    // Skip fill-driven removals (the FillsDelta carries the real fill).
414    if let Some(ref reason) = cancel.reason
415        && (reason == "full_fill" || reason == "partial_fill")
416    {
417        log::debug!(
418            "Skipping fill-driven cancel: order_id={}, reason={reason}",
419            cancel.order_id,
420        );
421        return;
422    }
423
424    let venue_order_id = VenueOrderId::new(&cancel.order_id);
425    let resolved_id = cancel
426        .cli_ord_id
427        .as_ref()
428        .map(|id| resolve_client_order_id(id, truncated_id_map))
429        .or_else(|| venue_client_map.load().get(&cancel.order_id).copied());
430
431    if let Some(client_order_id) = resolved_id
432        && let Some(identity) = state.lookup_identity(&client_order_id)
433    {
434        let ts_event = ts_init;
435        ensure_accepted_emitted(
436            client_order_id,
437            venue_order_id,
438            account_id,
439            &identity,
440            state,
441            emitter,
442            ts_event,
443            ts_init,
444        );
445        let canceled = OrderCanceled::new(
446            emitter.trader_id(),
447            identity.strategy_id,
448            identity.instrument_id,
449            client_order_id,
450            UUID4::new(),
451            ts_event,
452            ts_init,
453            false,
454            Some(venue_order_id),
455            Some(account_id),
456            cancel.reason.as_deref().map(Ustr::from),
457        );
458        emitter.send_order_event(OrderEventAny::Canceled(canceled));
459        state.cleanup_terminal(&client_order_id);
460        return;
461    }
462
463    // External fallback: build a status report from the side caches.
464    let Some(instrument_id) = order_instrument_map.load().get(&cancel.order_id).copied() else {
465        log::warn!(
466            "Cannot resolve instrument for cancel: order_id={}, \
467             order not seen in previous delta",
468            cancel.order_id
469        );
470        return;
471    };
472
473    let Some(quantity) = venue_order_qty.load().get(&cancel.order_id).copied() else {
474        log::warn!(
475            "Cannot resolve quantity for cancel: order_id={}, skipping",
476            cancel.order_id
477        );
478        return;
479    };
480
481    let report = OrderStatusReport::new(
482        account_id,
483        instrument_id,
484        resolved_id,
485        venue_order_id,
486        None,
487        OrderType::Limit,
488        TimeInForce::Gtc,
489        OrderStatus::Canceled,
490        quantity,
491        Quantity::zero(0),
492        ts_init,
493        ts_init,
494        ts_init,
495        None,
496    );
497    let report = if let Some(ref reason) = cancel.reason
498        && !reason.is_empty()
499    {
500        report.with_cancel_reason(reason.clone())
501    } else {
502        report
503    };
504    emitter.send_order_status_report(report);
505}
506
507/// Dispatches a Kraken Futures `FillsDelta` message.
508#[expect(clippy::too_many_arguments)]
509pub fn fills_delta(
510    fills_delta: &KrakenFuturesFillsDelta,
511    state: &WsDispatchState,
512    emitter: &ExecutionEventEmitter,
513    instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
514    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
515    venue_client_map: &Arc<AtomicMap<String, ClientOrderId>>,
516    account_id: AccountId,
517    ts_init: UnixNanos,
518) {
519    let instruments = instruments.load();
520    let venue_clients = venue_client_map.load();
521
522    for fill in &fills_delta.fills {
523        single_fill(
524            fill,
525            state,
526            emitter,
527            &instruments,
528            truncated_id_map,
529            &venue_clients,
530            account_id,
531            ts_init,
532        );
533    }
534}
535
536#[expect(clippy::too_many_arguments)]
537fn single_fill(
538    fill: &KrakenFuturesFill,
539    state: &WsDispatchState,
540    emitter: &ExecutionEventEmitter,
541    instruments: &AHashMap<InstrumentId, InstrumentAny>,
542    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
543    venue_client_map: &AHashMap<String, ClientOrderId>,
544    account_id: AccountId,
545    ts_init: UnixNanos,
546) {
547    let product_id = match &fill.instrument {
548        Some(id) => id.as_str(),
549        None => {
550            log::warn!("Fill missing instrument field: fill_id={}", fill.fill_id);
551            return;
552        }
553    };
554
555    let Some(instrument) = lookup_instrument_in_snapshot(instruments, product_id) else {
556        log::warn!("No instrument for product_id: {product_id}");
557        return;
558    };
559
560    let mut report = match parse_futures_ws_fill_report(fill, instrument, account_id, ts_init) {
561        Ok(report) => report,
562        Err(e) => {
563            log::error!("Failed to parse futures fill report: {e}");
564            return;
565        }
566    };
567
568    let resolved_id = fill
569        .cli_ord_id
570        .as_deref()
571        .filter(|s| !s.is_empty())
572        .map(|id| resolve_client_order_id(id, truncated_id_map))
573        .or_else(|| venue_client_map.get(&fill.order_id).copied());
574
575    if let Some(cid) = resolved_id
576        && state.filled_orders.contains(&cid)
577    {
578        log::debug!(
579            "Skipping stale fill for filled order: cid={cid}, order_id={}",
580            fill.order_id,
581        );
582        return;
583    }
584
585    if let Some(client_order_id) = resolved_id {
586        report.client_order_id = Some(client_order_id);
587
588        if let Some(identity) = state.lookup_identity(&client_order_id) {
589            if state.check_and_insert_trade(report.trade_id) {
590                log::debug!(
591                    "Skipping duplicate fill for {client_order_id}: trade_id={}",
592                    report.trade_id
593                );
594                return;
595            }
596            ensure_accepted_emitted(
597                client_order_id,
598                report.venue_order_id,
599                account_id,
600                &identity,
601                state,
602                emitter,
603                report.ts_event,
604                ts_init,
605            );
606            let filled = fill_report_to_order_filled(
607                &report,
608                emitter.trader_id(),
609                &identity,
610                instrument.quote_currency(),
611                client_order_id,
612            );
613            emitter.send_order_event(OrderEventAny::Filled(filled));
614
615            // Update cumulative filled and cleanup on terminal fill.
616            let previous = state
617                .previous_filled_qty(&client_order_id)
618                .unwrap_or_else(|| Quantity::zero(instrument.size_precision()));
619            let cumulative = previous + report.last_qty;
620            state.record_filled_qty(client_order_id, cumulative);
621
622            if cumulative >= identity.quantity {
623                state.insert_filled(client_order_id);
624                state.cleanup_terminal(&client_order_id);
625                return;
626            }
627
628            if let Some(pending) = state.pending_removal(&client_order_id)
629                && cumulative >= pending.venue_filled
630            {
631                state.remove_pending_removal(&client_order_id);
632
633                let canceled = OrderCanceled::new(
634                    emitter.trader_id(),
635                    identity.strategy_id,
636                    identity.instrument_id,
637                    client_order_id,
638                    UUID4::new(),
639                    pending.ts_event,
640                    ts_init,
641                    false,
642                    Some(report.venue_order_id),
643                    Some(account_id),
644                    pending.reason.as_deref().map(Ustr::from),
645                );
646                emitter.send_order_event(OrderEventAny::Canceled(canceled));
647                state.cleanup_terminal(&client_order_id);
648            }
649            return;
650        }
651    }
652
653    // External fallback.
654    if state.check_and_insert_trade(report.trade_id) {
655        log::debug!(
656            "Skipping duplicate external fill: trade_id={}",
657            report.trade_id
658        );
659        return;
660    }
661    emitter.send_fill_report(report);
662}
663
664#[inline]
665fn millis_to_nanos(millis: i64) -> UnixNanos {
666    UnixNanos::from((millis as u64) * 1_000_000)
667}