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