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};
35
36use super::{
37    DeltaSnapshot, OrderIdentity, WsDispatchState, ensure_accepted_emitted,
38    fill_report_to_order_filled, resolve_client_order_id,
39};
40use crate::{
41    common::lookup_instrument_in_snapshot,
42    websocket::futures::{
43        messages::{
44            KrakenFuturesFill, KrakenFuturesFillsDelta, KrakenFuturesOpenOrdersCancel,
45            KrakenFuturesOpenOrdersDelta,
46        },
47        parse::{parse_futures_ws_fill_report, parse_futures_ws_order_status_report},
48    },
49};
50
51/// Dispatches a Kraken Futures `OpenOrdersDelta` message.
52///
53/// Fill-driven cancel deltas (`is_cancel=true` with reason `full_fill` /
54/// `partial_fill`) are skipped - the corresponding `FillsDelta` carries the
55/// real fill, so emitting a synthetic Canceled here would race with the
56/// genuine `OrderFilled`.
57#[expect(clippy::too_many_arguments)]
58pub fn open_orders_delta(
59    delta: &KrakenFuturesOpenOrdersDelta,
60    state: &WsDispatchState,
61    emitter: &ExecutionEventEmitter,
62    instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
63    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
64    order_instrument_map: &Arc<AtomicMap<String, InstrumentId>>,
65    venue_client_map: &Arc<AtomicMap<String, ClientOrderId>>,
66    venue_order_qty: &Arc<AtomicMap<String, Quantity>>,
67    account_id: AccountId,
68    ts_init: UnixNanos,
69) {
70    if delta.is_fill_driven_cancel() {
71        log::debug!(
72            "Skipping fill-driven open_orders delta: order_id={}, reason={:?}",
73            delta.order.order_id,
74            delta.reason,
75        );
76        return;
77    }
78
79    let product_id = delta.order.instrument.as_str();
80    let instruments = instruments.load();
81    let Some(instrument) = lookup_instrument_in_snapshot(&instruments, product_id) else {
82        log::warn!("No instrument for product_id: {product_id}");
83        return;
84    };
85
86    // Cache instrument and qty by venue order id so cancel-only messages
87    // (which arrive without the order body) can be reconstructed for the
88    // external fallback path.
89    order_instrument_map.insert(delta.order.order_id.clone(), instrument.id());
90    let Ok(qty) = Quantity::from_decimal_dp(delta.order.qty, instrument.size_precision()) else {
91        log::error!("Failed to parse order quantity: {}", delta.order.qty);
92        return;
93    };
94    venue_order_qty.insert(delta.order.order_id.clone(), qty);
95
96    let resolved_id = delta
97        .order
98        .cli_ord_id
99        .as_ref()
100        .map(|id| resolve_client_order_id(id, truncated_id_map));
101
102    // Stale-report suppression: an order that already reached the filled
103    // terminal state should not produce more events even if a late delta
104    // arrives. `filled_orders` persists past `cleanup_terminal` precisely
105    // for this check.
106    if let Some(cid) = resolved_id
107        && state.filled_orders.contains(&cid)
108    {
109        log::debug!(
110            "Skipping stale open_orders delta for filled order: cid={cid}, order_id={}",
111            delta.order.order_id,
112        );
113        return;
114    }
115
116    if let Some(client_order_id) = resolved_id {
117        venue_client_map.insert(delta.order.order_id.clone(), client_order_id);
118
119        if let Some(identity) = state.lookup_identity(&client_order_id) {
120            delta_tracked(
121                delta,
122                client_order_id,
123                &identity,
124                instrument,
125                state,
126                emitter,
127                account_id,
128                ts_init,
129            );
130            return;
131        }
132    }
133
134    // External / untracked: fall back to a status report.
135    match parse_futures_ws_order_status_report(
136        &delta.order,
137        delta.is_cancel,
138        delta.reason.as_deref(),
139        instrument,
140        account_id,
141        ts_init,
142    ) {
143        Ok(mut report) => {
144            if let Some(cid) = resolved_id {
145                report = report.with_client_order_id(cid);
146            }
147            emitter.send_order_status_report(report);
148        }
149        Err(e) => log::error!("Failed to parse futures order status report: {e}"),
150    }
151}
152
153#[expect(clippy::too_many_arguments)]
154fn delta_tracked(
155    delta: &KrakenFuturesOpenOrdersDelta,
156    client_order_id: ClientOrderId,
157    identity: &OrderIdentity,
158    instrument: &InstrumentAny,
159    state: &WsDispatchState,
160    emitter: &ExecutionEventEmitter,
161    account_id: AccountId,
162    ts_init: UnixNanos,
163) {
164    let venue_order_id = VenueOrderId::new(&delta.order.order_id);
165    let ts_event = millis_to_nanos(delta.order.last_update_time);
166    let Ok(new_filled) = Quantity::from_decimal_dp(delta.order.filled, instrument.size_precision())
167    else {
168        log::error!("Failed to parse filled quantity: {}", delta.order.filled);
169        return;
170    };
171
172    if delta.is_cancel {
173        ensure_accepted_emitted(
174            client_order_id,
175            venue_order_id,
176            account_id,
177            identity,
178            state,
179            emitter,
180            ts_event,
181            ts_init,
182        );
183        let canceled = OrderCanceled::new(
184            emitter.trader_id(),
185            identity.strategy_id,
186            identity.instrument_id,
187            client_order_id,
188            UUID4::new(),
189            ts_event,
190            ts_init,
191            false,
192            Some(venue_order_id),
193            Some(account_id),
194        );
195        emitter.send_order_event(OrderEventAny::Canceled(canceled));
196        state.cleanup_terminal(&client_order_id);
197        return;
198    }
199
200    let already_accepted = state.emitted_accepted.contains(&client_order_id);
201    ensure_accepted_emitted(
202        client_order_id,
203        venue_order_id,
204        account_id,
205        identity,
206        state,
207        emitter,
208        ts_event,
209        ts_init,
210    );
211
212    let Ok(qty) = Quantity::from_decimal_dp(delta.order.qty, instrument.size_precision()) else {
213        log::error!("Failed to parse order quantity: {}", delta.order.qty);
214        return;
215    };
216    let snapshot = DeltaSnapshot::new(
217        qty,
218        new_filled,
219        delta.order.limit_price,
220        delta.order.stop_price,
221    );
222
223    if !already_accepted {
224        // First delta seen for this order: the placement Accepted is enough.
225        state.record_delta_snapshot(client_order_id, snapshot);
226        return;
227    }
228
229    // Follow-up delta. The two emission-relevant signals are independent:
230    //   * filled increased       -> partial-fill notification (FillsDelta has it,
231    //                               nothing to emit from here)
232    //   * non-fill field changed -> modify acknowledgement (emit OrderUpdated)
233    // Both can be true simultaneously when a user amends a partially filled
234    // order, so check the modify branch regardless of fill movement.
235    let previous = state.previous_delta_snapshot(&client_order_id);
236    state.record_delta_snapshot(client_order_id, snapshot);
237
238    let non_fill_changed = previous.is_some_and(|prev| !snapshot.non_fill_fields_match(&prev));
239    if !non_fill_changed {
240        return;
241    }
242
243    // Modify ack: refresh tracked quantity (size may have changed) and emit
244    // OrderUpdated so the engine clears PendingUpdate.
245    state.update_identity_quantity(&client_order_id, qty);
246    let price = match delta
247        .order
248        .limit_price
249        .map(|p| Price::from_decimal_dp(p, instrument.price_precision()))
250        .transpose()
251    {
252        Ok(price) => price,
253        Err(e) => {
254            log::error!("Failed to parse limit price: {e}");
255            return;
256        }
257    };
258    let trigger_price = match delta
259        .order
260        .stop_price
261        .map(|p| Price::from_decimal_dp(p, instrument.price_precision()))
262        .transpose()
263    {
264        Ok(price) => price,
265        Err(e) => {
266            log::error!("Failed to parse stop price: {e}");
267            return;
268        }
269    };
270
271    let updated = OrderUpdated::new(
272        emitter.trader_id(),
273        identity.strategy_id,
274        identity.instrument_id,
275        client_order_id,
276        qty,
277        UUID4::new(),
278        ts_event,
279        ts_init,
280        false,
281        Some(venue_order_id),
282        Some(account_id),
283        price,
284        trigger_price,
285        None,
286        false,
287    );
288    emitter.send_order_event(OrderEventAny::Updated(updated));
289}
290
291/// Dispatches a Kraken Futures `OpenOrdersCancel` (cancel-only) message.
292#[expect(clippy::too_many_arguments)]
293pub fn open_orders_cancel(
294    cancel: &KrakenFuturesOpenOrdersCancel,
295    state: &WsDispatchState,
296    emitter: &ExecutionEventEmitter,
297    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
298    order_instrument_map: &Arc<AtomicMap<String, InstrumentId>>,
299    venue_client_map: &Arc<AtomicMap<String, ClientOrderId>>,
300    venue_order_qty: &Arc<AtomicMap<String, Quantity>>,
301    account_id: AccountId,
302    ts_init: UnixNanos,
303) {
304    // Skip fill-driven removals (the FillsDelta carries the real fill).
305    if let Some(ref reason) = cancel.reason
306        && (reason == "full_fill" || reason == "partial_fill")
307    {
308        log::debug!(
309            "Skipping fill-driven cancel: order_id={}, reason={reason}",
310            cancel.order_id,
311        );
312        return;
313    }
314
315    let venue_order_id = VenueOrderId::new(&cancel.order_id);
316    let resolved_id = cancel
317        .cli_ord_id
318        .as_ref()
319        .map(|id| resolve_client_order_id(id, truncated_id_map))
320        .or_else(|| venue_client_map.load().get(&cancel.order_id).copied());
321
322    if let Some(client_order_id) = resolved_id
323        && let Some(identity) = state.lookup_identity(&client_order_id)
324    {
325        let ts_event = ts_init;
326        ensure_accepted_emitted(
327            client_order_id,
328            venue_order_id,
329            account_id,
330            &identity,
331            state,
332            emitter,
333            ts_event,
334            ts_init,
335        );
336        let canceled = OrderCanceled::new(
337            emitter.trader_id(),
338            identity.strategy_id,
339            identity.instrument_id,
340            client_order_id,
341            UUID4::new(),
342            ts_event,
343            ts_init,
344            false,
345            Some(venue_order_id),
346            Some(account_id),
347        );
348        emitter.send_order_event(OrderEventAny::Canceled(canceled));
349        state.cleanup_terminal(&client_order_id);
350        return;
351    }
352
353    // External fallback: build a status report from the side caches.
354    let Some(instrument_id) = order_instrument_map.load().get(&cancel.order_id).copied() else {
355        log::warn!(
356            "Cannot resolve instrument for cancel: order_id={}, \
357             order not seen in previous delta",
358            cancel.order_id
359        );
360        return;
361    };
362
363    let Some(quantity) = venue_order_qty.load().get(&cancel.order_id).copied() else {
364        log::warn!(
365            "Cannot resolve quantity for cancel: order_id={}, skipping",
366            cancel.order_id
367        );
368        return;
369    };
370
371    let report = OrderStatusReport::new(
372        account_id,
373        instrument_id,
374        resolved_id,
375        venue_order_id,
376        None,
377        OrderType::Limit,
378        TimeInForce::Gtc,
379        OrderStatus::Canceled,
380        quantity,
381        Quantity::zero(0),
382        ts_init,
383        ts_init,
384        ts_init,
385        None,
386    );
387    let report = if let Some(ref reason) = cancel.reason
388        && !reason.is_empty()
389    {
390        report.with_cancel_reason(reason.clone())
391    } else {
392        report
393    };
394    emitter.send_order_status_report(report);
395}
396
397/// Dispatches a Kraken Futures `FillsDelta` message.
398#[expect(clippy::too_many_arguments)]
399pub fn fills_delta(
400    fills_delta: &KrakenFuturesFillsDelta,
401    state: &WsDispatchState,
402    emitter: &ExecutionEventEmitter,
403    instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
404    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
405    venue_client_map: &Arc<AtomicMap<String, ClientOrderId>>,
406    account_id: AccountId,
407    ts_init: UnixNanos,
408) {
409    let instruments = instruments.load();
410    let venue_clients = venue_client_map.load();
411
412    for fill in &fills_delta.fills {
413        single_fill(
414            fill,
415            state,
416            emitter,
417            &instruments,
418            truncated_id_map,
419            &venue_clients,
420            account_id,
421            ts_init,
422        );
423    }
424}
425
426#[expect(clippy::too_many_arguments)]
427fn single_fill(
428    fill: &KrakenFuturesFill,
429    state: &WsDispatchState,
430    emitter: &ExecutionEventEmitter,
431    instruments: &AHashMap<InstrumentId, InstrumentAny>,
432    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
433    venue_client_map: &AHashMap<String, ClientOrderId>,
434    account_id: AccountId,
435    ts_init: UnixNanos,
436) {
437    let product_id = match &fill.instrument {
438        Some(id) => id.as_str(),
439        None => {
440            log::warn!("Fill missing instrument field: fill_id={}", fill.fill_id);
441            return;
442        }
443    };
444
445    let Some(instrument) = lookup_instrument_in_snapshot(instruments, product_id) else {
446        log::warn!("No instrument for product_id: {product_id}");
447        return;
448    };
449
450    let mut report = match parse_futures_ws_fill_report(fill, instrument, account_id, ts_init) {
451        Ok(report) => report,
452        Err(e) => {
453            log::error!("Failed to parse futures fill report: {e}");
454            return;
455        }
456    };
457
458    let resolved_id = fill
459        .cli_ord_id
460        .as_deref()
461        .filter(|s| !s.is_empty())
462        .map(|id| resolve_client_order_id(id, truncated_id_map))
463        .or_else(|| venue_client_map.get(&fill.order_id).copied());
464
465    if let Some(cid) = resolved_id
466        && state.filled_orders.contains(&cid)
467    {
468        log::debug!(
469            "Skipping stale fill for filled order: cid={cid}, order_id={}",
470            fill.order_id,
471        );
472        return;
473    }
474
475    if let Some(client_order_id) = resolved_id {
476        report.client_order_id = Some(client_order_id);
477
478        if let Some(identity) = state.lookup_identity(&client_order_id) {
479            if state.check_and_insert_trade(report.trade_id) {
480                log::debug!(
481                    "Skipping duplicate fill for {client_order_id}: trade_id={}",
482                    report.trade_id
483                );
484                return;
485            }
486            ensure_accepted_emitted(
487                client_order_id,
488                report.venue_order_id,
489                account_id,
490                &identity,
491                state,
492                emitter,
493                report.ts_event,
494                ts_init,
495            );
496            let filled = fill_report_to_order_filled(
497                &report,
498                emitter.trader_id(),
499                &identity,
500                instrument.quote_currency(),
501                client_order_id,
502            );
503            emitter.send_order_event(OrderEventAny::Filled(filled));
504
505            // Update cumulative filled and cleanup on terminal fill.
506            let previous = state
507                .previous_filled_qty(&client_order_id)
508                .unwrap_or_else(|| Quantity::zero(instrument.size_precision()));
509            let cumulative = previous + report.last_qty;
510            state.record_filled_qty(client_order_id, cumulative);
511
512            if cumulative >= identity.quantity {
513                state.insert_filled(client_order_id);
514                state.cleanup_terminal(&client_order_id);
515            }
516            return;
517        }
518    }
519
520    // External fallback.
521    if state.check_and_insert_trade(report.trade_id) {
522        log::debug!(
523            "Skipping duplicate external fill: trade_id={}",
524            report.trade_id
525        );
526        return;
527    }
528    emitter.send_fill_report(report);
529}
530
531#[inline]
532fn millis_to_nanos(millis: i64) -> UnixNanos {
533    UnixNanos::from((millis as u64) * 1_000_000)
534}