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};
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#[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 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 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 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 state.record_delta_snapshot(client_order_id, snapshot);
226 return;
227 }
228
229 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 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#[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 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 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#[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 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 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}