Skip to main content

nautilus_kraken/websocket/dispatch/
spot.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 Spot v2 API.
17//!
18//! A single spot execution can carry both a status update (handled via the
19//! order event path) and a fill (when `exec_id` is present, handled via the
20//! fill path). Tracked orders emit typed events; external orders fall through
21//! to reports.
22
23use std::sync::Arc;
24
25use nautilus_core::{AtomicMap, UUID4, UnixNanos};
26use nautilus_live::ExecutionEventEmitter;
27use nautilus_model::{
28    enums::OrderStatus,
29    events::{
30        OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired, OrderTriggered, OrderUpdated,
31    },
32    identifiers::{AccountId, ClientOrderId, InstrumentId},
33    instruments::{Instrument, InstrumentAny},
34    reports::{FillReport, OrderStatusReport},
35    types::Quantity,
36};
37use rust_decimal::Decimal;
38
39use super::{
40    OrderIdentity, WsDispatchState, ensure_accepted_emitted, fill_report_to_order_filled,
41    resolve_client_order_id,
42};
43use crate::{
44    common::lookup_instrument_in_snapshot,
45    websocket::spot_v2::{
46        enums::KrakenExecType,
47        messages::KrakenWsExecutionData,
48        parse::{parse_ws_fill_report, parse_ws_order_status_report},
49    },
50};
51
52/// Dispatches a Kraken Spot v2 execution message.
53#[expect(clippy::too_many_arguments)]
54pub fn execution(
55    exec: &KrakenWsExecutionData,
56    state: &WsDispatchState,
57    emitter: &ExecutionEventEmitter,
58    instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
59    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
60    order_qty_cache: &Arc<AtomicMap<String, Decimal>>,
61    account_id: AccountId,
62    ts_init: UnixNanos,
63) {
64    execution_inner(
65        exec,
66        state,
67        emitter,
68        instruments,
69        truncated_id_map,
70        order_qty_cache,
71        account_id,
72        ts_init,
73    );
74
75    // Run terminal cache cleanup regardless of which early return the inner
76    // dispatch hit (symbol miss, instrument miss, stale-filled suppression,
77    // parse error). Keying eviction off `exec.order_id` means it does not
78    // depend on `cl_ord_id` or identity resolution succeeding.
79    if is_terminal_exec_type(exec.exec_type) {
80        state.forget_order_symbol(&exec.order_id);
81        state.forget_order_client_id(&exec.order_id);
82    }
83}
84
85#[expect(clippy::too_many_arguments)]
86fn execution_inner(
87    exec: &KrakenWsExecutionData,
88    state: &WsDispatchState,
89    emitter: &ExecutionEventEmitter,
90    instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
91    truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
92    order_qty_cache: &Arc<AtomicMap<String, Decimal>>,
93    account_id: AccountId,
94    ts_init: UnixNanos,
95) {
96    // Resolve the trading symbol. Per Kraken's executions docs, follow-up
97    // frames (`new`, `amended`, `restated`, `status`) carry only changed
98    // fields and omit `symbol`. We cache the symbol from the first frame for
99    // the venue order id and consult it here when the current frame omits it.
100    let cached_symbol;
101    let symbol = match exec.symbol.as_deref() {
102        Some(s) => {
103            state.cache_order_symbol(&exec.order_id, s);
104            s
105        }
106        None => match state.lookup_order_symbol(&exec.order_id) {
107            Some(s) => {
108                cached_symbol = s;
109                cached_symbol.as_str()
110            }
111            None => {
112                log::debug!(
113                    "Execution message without symbol and no cached mapping: \
114                     exec_type={:?}, order_id={}",
115                    exec.exec_type,
116                    exec.order_id
117                );
118                return;
119            }
120        },
121    };
122    let instruments = instruments.load();
123    let Some(instrument) = lookup_instrument_in_snapshot(&instruments, symbol) else {
124        log::warn!("No instrument for symbol: {symbol}");
125        return;
126    };
127
128    // Mirror the existing behaviour: cache the order quantity by truncated cli
129    // ord id so the parser can fall back to it for quote-quantity orders.
130    let cached_qty = exec
131        .cl_ord_id
132        .as_ref()
133        .and_then(|id| order_qty_cache.load().get(id).copied());
134    if let (Some(qty), Some(cl_ord_id)) = (exec.order_qty, &exec.cl_ord_id) {
135        order_qty_cache.insert(cl_ord_id.clone(), qty);
136    }
137
138    // Resolve the `ClientOrderId`. When the frame carries `cl_ord_id` we
139    // resolve it through the truncation map and seed the venue-id cache so
140    // later delta frames (which routinely omit `cl_ord_id`) can recover it.
141    // When `cl_ord_id` is absent we consult the cache; without this lookup a
142    // tracked order's delta `new` would fall through to the untracked report
143    // path and the strategy would never see `OrderAccepted` (issue #4051).
144    let resolved_id = match exec.cl_ord_id.as_ref() {
145        Some(id) => {
146            let cid = resolve_client_order_id(id, truncated_id_map);
147            state.cache_order_client_id(&exec.order_id, cid);
148            Some(cid)
149        }
150        None => state.lookup_order_client_id(&exec.order_id),
151    };
152
153    // Stale-report suppression for previously-tracked orders that already
154    // reached the filled terminal state.
155    if let Some(cid) = resolved_id
156        && state.filled_orders.contains(&cid)
157    {
158        log::debug!(
159            "Skipping stale spot execution for filled order: cid={cid}, order_id={}",
160            exec.order_id,
161        );
162        return;
163    }
164
165    let identity = resolved_id.and_then(|cid| state.lookup_identity(&cid));
166
167    // Status update.
168    match parse_ws_order_status_report(exec, instrument, account_id, cached_qty, ts_init) {
169        Ok(mut report) => {
170            if let Some(cid) = resolved_id {
171                report = report.with_client_order_id(cid);
172            }
173
174            if let (Some(client_order_id), Some(identity)) = (resolved_id, identity.as_ref()) {
175                status_tracked(
176                    &report,
177                    exec.exec_type,
178                    exec.exec_id.is_some(),
179                    client_order_id,
180                    identity,
181                    state,
182                    emitter,
183                    account_id,
184                    ts_init,
185                );
186            } else {
187                emitter.send_order_status_report(report);
188            }
189        }
190        Err(e) => log::error!("Failed to parse order status report: {e}"),
191    }
192
193    // Fill (when present).
194    if exec.exec_id.is_some() {
195        match parse_ws_fill_report(exec, instrument, account_id, ts_init) {
196            Ok(mut report) => {
197                if let Some(cid) = resolved_id {
198                    report.client_order_id = Some(cid);
199                }
200
201                if let (Some(client_order_id), Some(identity)) = (resolved_id, identity.as_ref()) {
202                    fill_tracked(
203                        &report,
204                        client_order_id,
205                        identity,
206                        instrument,
207                        state,
208                        emitter,
209                        account_id,
210                        ts_init,
211                    );
212                } else {
213                    if state.check_and_insert_trade(report.trade_id) {
214                        log::debug!(
215                            "Skipping duplicate external spot fill: trade_id={}",
216                            report.trade_id
217                        );
218                        return;
219                    }
220                    emitter.send_fill_report(report);
221                }
222            }
223            Err(e) => log::error!("Failed to parse fill report: {e}"),
224        }
225    }
226}
227
228#[expect(clippy::too_many_arguments)]
229fn status_tracked(
230    report: &OrderStatusReport,
231    exec_type: KrakenExecType,
232    has_fill: bool,
233    client_order_id: ClientOrderId,
234    identity: &OrderIdentity,
235    state: &WsDispatchState,
236    emitter: &ExecutionEventEmitter,
237    account_id: AccountId,
238    ts_init: UnixNanos,
239) {
240    let venue_order_id = report.venue_order_id;
241    let ts_event = report.ts_last;
242    let trader_id = emitter.trader_id();
243
244    // Amended (user modify) and Restated (engine adjustment) both surface
245    // post-modify state. Refresh tracked quantity (size may have changed) and
246    // emit OrderUpdated so the engine clears PendingUpdate.
247    if matches!(
248        exec_type,
249        KrakenExecType::Amended | KrakenExecType::Restated
250    ) && state.emitted_accepted.contains(&client_order_id)
251    {
252        state.update_identity_quantity(&client_order_id, report.quantity);
253        let updated = OrderUpdated::new(
254            trader_id,
255            identity.strategy_id,
256            identity.instrument_id,
257            client_order_id,
258            report.quantity,
259            UUID4::new(),
260            ts_event,
261            ts_init,
262            false,
263            Some(venue_order_id),
264            Some(account_id),
265            report.price,
266            report.trigger_price,
267            None,
268            false,
269        );
270        emitter.send_order_event(OrderEventAny::Updated(updated));
271        return;
272    }
273
274    match report.order_status {
275        OrderStatus::Accepted => {
276            if !state.insert_accepted(client_order_id) {
277                // Already accepted; this is a redundant New / Restated / Status
278                // exec. The strategy already saw OrderAccepted; nothing to emit.
279                return;
280            }
281            let accepted = OrderAccepted::new(
282                trader_id,
283                identity.strategy_id,
284                identity.instrument_id,
285                client_order_id,
286                venue_order_id,
287                account_id,
288                UUID4::new(),
289                ts_event,
290                ts_init,
291                false,
292            );
293            emitter.send_order_event(OrderEventAny::Accepted(accepted));
294        }
295        OrderStatus::Triggered => {
296            // Stop / take-profit transition. Synthesize Accepted first if the
297            // venue compressed placement and trigger into one message.
298            ensure_accepted_emitted(
299                client_order_id,
300                venue_order_id,
301                account_id,
302                identity,
303                state,
304                emitter,
305                ts_event,
306                ts_init,
307            );
308            let triggered = OrderTriggered::new(
309                trader_id,
310                identity.strategy_id,
311                identity.instrument_id,
312                client_order_id,
313                UUID4::new(),
314                ts_event,
315                ts_init,
316                false,
317                Some(venue_order_id),
318                Some(account_id),
319            );
320            emitter.send_order_event(OrderEventAny::Triggered(triggered));
321        }
322        OrderStatus::PartiallyFilled => {
323            // The fill itself is emitted from the trade-side of dispatch via
324            // fill_tracked; nothing to do here.
325        }
326
327        // Fill dispatch handles cleanup when the report includes a fill
328        OrderStatus::Filled if !has_fill => {
329            state.insert_filled(client_order_id);
330            state.cleanup_terminal(&client_order_id);
331        }
332        OrderStatus::Canceled => {
333            ensure_accepted_emitted(
334                client_order_id,
335                venue_order_id,
336                account_id,
337                identity,
338                state,
339                emitter,
340                ts_event,
341                ts_init,
342            );
343            let canceled = OrderCanceled::new(
344                trader_id,
345                identity.strategy_id,
346                identity.instrument_id,
347                client_order_id,
348                UUID4::new(),
349                ts_event,
350                ts_init,
351                false,
352                Some(venue_order_id),
353                Some(account_id),
354            );
355            emitter.send_order_event(OrderEventAny::Canceled(canceled));
356            state.cleanup_terminal(&client_order_id);
357        }
358        OrderStatus::Expired => {
359            ensure_accepted_emitted(
360                client_order_id,
361                venue_order_id,
362                account_id,
363                identity,
364                state,
365                emitter,
366                ts_event,
367                ts_init,
368            );
369            let expired = OrderExpired::new(
370                trader_id,
371                identity.strategy_id,
372                identity.instrument_id,
373                client_order_id,
374                UUID4::new(),
375                ts_event,
376                ts_init,
377                false,
378                Some(venue_order_id),
379                Some(account_id),
380            );
381            emitter.send_order_event(OrderEventAny::Expired(expired));
382            state.cleanup_terminal(&client_order_id);
383        }
384        _ => {}
385    }
386}
387
388#[expect(clippy::too_many_arguments)]
389fn fill_tracked(
390    report: &FillReport,
391    client_order_id: ClientOrderId,
392    identity: &OrderIdentity,
393    instrument: &InstrumentAny,
394    state: &WsDispatchState,
395    emitter: &ExecutionEventEmitter,
396    account_id: AccountId,
397    ts_init: UnixNanos,
398) {
399    if state.check_and_insert_trade(report.trade_id) {
400        log::debug!(
401            "Skipping duplicate spot fill for {client_order_id}: trade_id={}",
402            report.trade_id
403        );
404        return;
405    }
406
407    ensure_accepted_emitted(
408        client_order_id,
409        report.venue_order_id,
410        account_id,
411        identity,
412        state,
413        emitter,
414        report.ts_event,
415        ts_init,
416    );
417
418    let filled = fill_report_to_order_filled(
419        report,
420        emitter.trader_id(),
421        identity,
422        instrument.quote_currency(),
423        client_order_id,
424    );
425    emitter.send_order_event(OrderEventAny::Filled(filled));
426
427    let previous = state
428        .previous_filled_qty(&client_order_id)
429        .unwrap_or_else(|| Quantity::zero(instrument.size_precision()));
430    let cumulative = previous + report.last_qty;
431    state.record_filled_qty(client_order_id, cumulative);
432
433    if cumulative >= identity.quantity {
434        state.insert_filled(client_order_id);
435        state.cleanup_terminal(&client_order_id);
436    }
437}
438
439/// Returns true when this spot execution carries a terminal status that
440/// should remove the order from dispatch state.
441#[must_use]
442pub fn is_terminal_exec_type(exec_type: KrakenExecType) -> bool {
443    matches!(
444        exec_type,
445        KrakenExecType::Filled | KrakenExecType::Canceled | KrakenExecType::Expired
446    )
447}
448
449#[cfg(test)]
450mod tests {
451    use rstest::rstest;
452
453    use super::*;
454
455    #[rstest]
456    #[case::filled(KrakenExecType::Filled, true)]
457    #[case::canceled(KrakenExecType::Canceled, true)]
458    #[case::expired(KrakenExecType::Expired, true)]
459    #[case::new(KrakenExecType::New, false)]
460    #[case::trade(KrakenExecType::Trade, false)]
461    #[case::pending_new(KrakenExecType::PendingNew, false)]
462    fn test_is_terminal_exec_type(#[case] exec_type: KrakenExecType, #[case] expected: bool) {
463        assert_eq!(is_terminal_exec_type(exec_type), expected);
464    }
465}