1use 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#[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 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 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 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 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 state.record_delta_snapshot(client_order_id, snapshot);
258 return;
259 }
260
261 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 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#[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#[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 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 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#[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 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 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}