1use std::{
68 collections::VecDeque,
69 sync::atomic::{AtomicBool, Ordering},
70};
71
72use dashmap::{DashMap, DashSet};
73use nautilus_common::cache::fifo::FifoCache;
74use nautilus_core::{UUID4, UnixNanos};
75use nautilus_live::{ExecutionEventEmitter, execution::context::OrderContext};
76use nautilus_model::{
77 enums::{OrderStatus, OrderType},
78 events::{
79 OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired, OrderFilled, OrderRejected,
80 OrderTriggered, OrderUpdated,
81 },
82 identifiers::{AccountId, ClientOrderId, TradeId, VenueOrderId},
83 reports::{FillReport, OrderStatusReport},
84 types::{Price, Quantity},
85};
86use parking_lot::Mutex;
87use ustr::Ustr;
88
89use crate::{
90 common::consts::HYPERLIQUID_POST_ONLY_WOULD_MATCH,
91 http::models::HyperliquidExchangePlaceOrderRequest,
92};
93
94pub const DEDUP_CAPACITY: usize = 10_000;
95
96pub const MAX_PENDING_MODIFY_INTENTS: usize = 32;
100
101#[derive(Debug, Clone)]
110pub struct ModifyIntent {
111 pub generation: u64,
113 pub old_venue_order_id: Option<VenueOrderId>,
115 pub target_qty: Quantity,
117 pub sent_request: Option<HyperliquidExchangePlaceOrderRequest>,
119}
120
121#[derive(Debug, Default)]
127struct ModifyChain {
128 intents: VecDeque<ModifyIntent>,
129 next_generation: u64,
130}
131
132#[derive(Debug)]
142pub struct WsDispatchState {
143 pub order_contexts: DashMap<ClientOrderId, OrderContext>,
152 pub emitted_accepted: DashSet<ClientOrderId>,
154 pending_submissions: DashSet<ClientOrderId>,
156 pending_submission_rejections: DashMap<ClientOrderId, OrderStatusReport>,
159 pub filled_orders: DashSet<ClientOrderId>,
164 emitted_trades: Mutex<FifoCache<TradeId, DEDUP_CAPACITY>>,
169 terminal_cloids: Mutex<FifoCache<Ustr, DEDUP_CAPACITY>>,
172 pub cached_venue_order_ids: DashMap<ClientOrderId, VenueOrderId>,
179 pending_modify_chains: DashMap<ClientOrderId, ModifyChain>,
191 pub buffered_fills: DashMap<ClientOrderId, Vec<FillReport>>,
196 pub order_filled_qty: DashMap<ClientOrderId, Quantity>,
199 pub pending_corrective: DashMap<ClientOrderId, (u64, HyperliquidExchangePlaceOrderRequest)>,
202 clearing: AtomicBool,
203}
204
205impl Default for WsDispatchState {
206 fn default() -> Self {
207 Self {
208 order_contexts: DashMap::new(),
209 emitted_accepted: DashSet::default(),
210 pending_submissions: DashSet::default(),
211 pending_submission_rejections: DashMap::new(),
212 filled_orders: DashSet::default(),
213 emitted_trades: Mutex::new(FifoCache::new()),
214 terminal_cloids: Mutex::new(FifoCache::new()),
215 cached_venue_order_ids: DashMap::new(),
216 pending_modify_chains: DashMap::new(),
217 buffered_fills: DashMap::new(),
218 order_filled_qty: DashMap::new(),
219 pending_corrective: DashMap::new(),
220 clearing: AtomicBool::new(false),
221 }
222 }
223}
224
225impl WsDispatchState {
226 #[must_use]
228 pub fn new() -> Self {
229 Self::default()
230 }
231
232 pub fn register_context(&self, context: OrderContext) {
235 self.order_contexts
236 .insert(context.identity.client_order_id, context);
237 }
238
239 #[must_use]
241 pub fn lookup_context(&self, client_order_id: &ClientOrderId) -> Option<OrderContext> {
242 self.order_contexts.get(client_order_id).map(|r| *r)
243 }
244
245 pub fn mark_submission_pending(&self, client_order_id: ClientOrderId) {
247 self.pending_submissions.insert(client_order_id);
248 }
249
250 #[must_use]
252 pub fn submission_pending(&self, client_order_id: &ClientOrderId) -> bool {
253 self.pending_submissions.contains(client_order_id)
254 }
255
256 pub fn buffer_submission_rejection(
258 &self,
259 client_order_id: ClientOrderId,
260 report: OrderStatusReport,
261 ) {
262 self.pending_submission_rejections
263 .insert(client_order_id, report);
264 }
265
266 #[must_use]
268 pub fn resolve_submission(&self, client_order_id: &ClientOrderId) -> Option<OrderStatusReport> {
269 self.pending_submissions.remove(client_order_id);
270 self.pending_submission_rejections
271 .remove(client_order_id)
272 .map(|(_, report)| report)
273 }
274
275 pub fn update_context_price(&self, client_order_id: &ClientOrderId, price: Option<Price>) {
278 if let Some(price) = price
279 && let Some(mut entry) = self.order_contexts.get_mut(client_order_id)
280 {
281 entry.price = Some(price);
282 }
283 }
284
285 pub fn update_context_quantity(&self, client_order_id: &ClientOrderId, quantity: Quantity) {
287 if let Some(mut entry) = self.order_contexts.get_mut(client_order_id) {
288 entry.quantity = quantity;
289 }
290 }
291
292 pub fn insert_accepted(&self, cid: ClientOrderId) {
294 self.evict_if_full(&self.emitted_accepted);
295 self.emitted_accepted.insert(cid);
296 }
297
298 pub fn insert_filled(&self, cid: ClientOrderId) -> bool {
303 self.evict_if_full(&self.filled_orders);
304 self.filled_orders.insert(cid)
305 }
306
307 pub fn check_and_insert_trade(&self, trade_id: TradeId) -> bool {
312 let mut set = self.emitted_trades.lock();
313 !set.insert(trade_id)
314 }
315
316 pub fn insert_terminal_cloid(&self, cloid: Ustr) {
323 let mut set = self.terminal_cloids.lock();
324 let _ = set.insert(cloid);
325 }
326
327 #[must_use]
330 pub fn terminal_cloid_seen(&self, cloid: &Ustr) -> bool {
331 let set = self.terminal_cloids.lock();
332 set.contains(cloid)
333 }
334
335 pub fn record_venue_order_id(
337 &self,
338 client_order_id: ClientOrderId,
339 venue_order_id: VenueOrderId,
340 ) {
341 self.cached_venue_order_ids
342 .insert(client_order_id, venue_order_id);
343 }
344
345 #[must_use]
347 pub fn cached_venue_order_id(&self, client_order_id: &ClientOrderId) -> Option<VenueOrderId> {
348 self.cached_venue_order_ids.get(client_order_id).map(|r| *r)
349 }
350
351 pub fn mark_pending_modify(
360 &self,
361 client_order_id: ClientOrderId,
362 old_venue_order_id: VenueOrderId,
363 target_qty: Quantity,
364 ) -> u64 {
365 let mut chain = self
366 .pending_modify_chains
367 .entry(client_order_id)
368 .or_default();
369 let generation = chain.next_generation;
370 chain.next_generation += 1;
371 chain.intents.push_back(ModifyIntent {
372 generation,
373 old_venue_order_id: Some(old_venue_order_id),
374 target_qty,
375 sent_request: None,
376 });
377
378 if chain.intents.len() > MAX_PENDING_MODIFY_INTENTS {
379 chain.intents.pop_front();
380 log::warn!(
381 "Modify chain for {client_order_id} exceeded {MAX_PENDING_MODIFY_INTENTS}; \
382 evicting oldest intent",
383 );
384 }
385 generation
386 }
387
388 pub fn clear_pending_modify(&self, client_order_id: &ClientOrderId) {
390 self.pending_modify_chains.remove(client_order_id);
391 }
392
393 pub fn clear_modify_generation(&self, client_order_id: &ClientOrderId, generation: u64) {
402 let Some(mut chain) = self.pending_modify_chains.get_mut(client_order_id) else {
403 return;
404 };
405 let removed_front_old = chain
406 .intents
407 .front()
408 .filter(|front| front.generation == generation)
409 .and_then(|front| front.old_venue_order_id);
410 chain
411 .intents
412 .retain(|intent| intent.generation != generation);
413
414 if let Some(old) = removed_front_old
415 && let Some(new_front) = chain.intents.front_mut()
416 {
417 new_front.old_venue_order_id = Some(old);
418 }
419 drop(chain);
420 self.pending_modify_chains
423 .remove_if(client_order_id, |_, chain| chain.intents.is_empty());
424 }
425
426 pub fn stash_modify_request(
429 &self,
430 client_order_id: ClientOrderId,
431 request: HyperliquidExchangePlaceOrderRequest,
432 ) {
433 if let Some(mut chain) = self.pending_modify_chains.get_mut(&client_order_id)
434 && let Some(back) = chain.intents.back_mut()
435 {
436 back.sent_request = Some(request);
437 } else {
438 log::debug!(
439 "Stash modify request for {client_order_id} with no pending intent; ignoring"
440 );
441 }
442 }
443
444 #[must_use]
446 pub fn modify_request(
447 &self,
448 client_order_id: &ClientOrderId,
449 ) -> Option<HyperliquidExchangePlaceOrderRequest> {
450 self.pending_modify_chains
451 .get(client_order_id)
452 .and_then(|chain| chain.intents.front().and_then(|i| i.sent_request.clone()))
453 }
454
455 pub fn claim_front_modify(
463 &self,
464 client_order_id: &ClientOrderId,
465 new_venue_order_id: VenueOrderId,
466 ) -> Option<ModifyIntent> {
467 let mut chain = self.pending_modify_chains.get_mut(client_order_id)?;
468 let claimed = chain.intents.pop_front();
469 if let Some(next) = chain.intents.front_mut() {
470 next.old_venue_order_id = Some(new_venue_order_id);
471 }
472 drop(chain);
473 self.pending_modify_chains
476 .remove_if(client_order_id, |_, chain| chain.intents.is_empty());
477 claimed
478 }
479
480 pub fn queue_corrective(
482 &self,
483 client_order_id: ClientOrderId,
484 oid: u64,
485 request: HyperliquidExchangePlaceOrderRequest,
486 ) {
487 self.pending_corrective
488 .insert(client_order_id, (oid, request));
489 }
490
491 #[must_use]
493 pub fn take_corrective(
494 &self,
495 client_order_id: &ClientOrderId,
496 ) -> Option<(u64, HyperliquidExchangePlaceOrderRequest)> {
497 self.pending_corrective
498 .remove(client_order_id)
499 .map(|(_, v)| v)
500 }
501
502 #[must_use]
504 pub fn has_pending_modify(&self, client_order_id: &ClientOrderId) -> bool {
505 self.pending_modify_chains
506 .get(client_order_id)
507 .is_some_and(|chain| !chain.intents.is_empty())
508 }
509
510 #[must_use]
512 pub fn pending_modify(&self, client_order_id: &ClientOrderId) -> Option<VenueOrderId> {
513 self.pending_modify_chains
514 .get(client_order_id)
515 .and_then(|chain| chain.intents.front().and_then(|i| i.old_venue_order_id))
516 }
517
518 #[must_use]
523 pub fn pending_modify_contains_old(
524 &self,
525 client_order_id: &ClientOrderId,
526 venue_order_id: VenueOrderId,
527 ) -> bool {
528 self.pending_modify_chains
529 .get(client_order_id)
530 .is_some_and(|chain| {
531 chain
532 .intents
533 .iter()
534 .any(|i| i.old_venue_order_id == Some(venue_order_id))
535 })
536 }
537
538 #[must_use]
540 pub fn pending_modify_target_qty(&self, client_order_id: &ClientOrderId) -> Option<Quantity> {
541 self.pending_modify_chains
542 .get(client_order_id)
543 .and_then(|chain| chain.intents.front().map(|i| i.target_qty))
544 }
545
546 pub fn buffer_fill(&self, client_order_id: ClientOrderId, fill: FillReport) {
548 self.buffered_fills
549 .entry(client_order_id)
550 .or_default()
551 .push(fill);
552 }
553
554 #[must_use]
556 pub fn drain_buffered_fills(&self, client_order_id: &ClientOrderId) -> Vec<FillReport> {
557 self.buffered_fills
558 .remove(client_order_id)
559 .map(|(_, v)| v)
560 .unwrap_or_default()
561 }
562
563 #[must_use]
565 pub fn buffered_fill_count(&self, client_order_id: &ClientOrderId) -> usize {
566 self.buffered_fills
567 .get(client_order_id)
568 .map_or(0, |r| r.len())
569 }
570
571 pub fn record_filled_qty(&self, client_order_id: ClientOrderId, qty: Quantity) {
573 self.order_filled_qty.insert(client_order_id, qty);
574 }
575
576 #[must_use]
578 pub fn previous_filled_qty(&self, client_order_id: &ClientOrderId) -> Option<Quantity> {
579 self.order_filled_qty.get(client_order_id).map(|r| *r)
580 }
581
582 pub fn cleanup_terminal(&self, client_order_id: &ClientOrderId) {
587 self.order_contexts.remove(client_order_id);
588 self.emitted_accepted.remove(client_order_id);
589 self.pending_submissions.remove(client_order_id);
590 self.pending_submission_rejections.remove(client_order_id);
591 self.cached_venue_order_ids.remove(client_order_id);
592 self.pending_modify_chains.remove(client_order_id);
593 self.pending_corrective.remove(client_order_id);
594 self.buffered_fills.remove(client_order_id);
595 self.order_filled_qty.remove(client_order_id);
596 }
597
598 fn evict_if_full(&self, set: &DashSet<ClientOrderId>) {
599 if set.len() >= DEDUP_CAPACITY
600 && self
601 .clearing
602 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed)
603 .is_ok()
604 {
605 set.clear();
606 self.clearing.store(false, Ordering::Release);
607 }
608 }
609}
610
611#[derive(Debug, Clone, Copy, PartialEq, Eq)]
613pub enum DispatchOutcome {
614 Tracked,
618 External,
623 Skip,
627}
628
629pub fn dispatch_order_event(
639 report: &OrderStatusReport,
640 state: &WsDispatchState,
641 emitter: &ExecutionEventEmitter,
642 ts_init: UnixNanos,
643) -> DispatchOutcome {
644 let Some(client_order_id) = report.client_order_id else {
645 return DispatchOutcome::External;
646 };
647
648 if state.filled_orders.contains(&client_order_id) {
649 log::debug!(
650 "Skipping stale report for filled order: cid={client_order_id}, status={:?}",
651 report.order_status,
652 );
653 return DispatchOutcome::Skip;
654 }
655
656 let client_order_id_str = client_order_id.as_str();
657 if client_order_id_str.starts_with("0x")
658 && state.terminal_cloid_seen(&Ustr::from(client_order_id_str))
659 {
660 log::debug!(
661 "Skipping stale terminal report for raw cloid: cid={client_order_id}, status={:?}",
662 report.order_status,
663 );
664 return DispatchOutcome::Skip;
665 }
666
667 let Some(context) = state.lookup_context(&client_order_id) else {
668 return DispatchOutcome::External;
669 };
670
671 match report.order_status {
672 OrderStatus::Accepted => {
673 handle_accepted(report, client_order_id, &context, state, emitter, ts_init)
674 }
675 OrderStatus::Triggered => {
676 handle_triggered(report, client_order_id, &context, state, emitter, ts_init)
677 }
678 OrderStatus::Canceled => {
679 handle_canceled(report, client_order_id, &context, state, emitter, ts_init)
680 }
681 OrderStatus::Expired => {
682 handle_expired(report, client_order_id, &context, state, emitter, ts_init)
683 }
684 OrderStatus::Rejected => {
685 handle_rejected(report, client_order_id, &context, state, emitter, ts_init)
686 }
687 OrderStatus::Filled => handle_filled_marker(client_order_id, state),
688 OrderStatus::PartiallyFilled => {
689 DispatchOutcome::Tracked
691 }
692 OrderStatus::PendingUpdate
693 | OrderStatus::PendingCancel
694 | OrderStatus::Submitted
695 | OrderStatus::Initialized
696 | OrderStatus::Denied
697 | OrderStatus::Released
698 | OrderStatus::Emulated
699 | OrderStatus::Voided => DispatchOutcome::Tracked,
700 }
701}
702
703pub fn dispatch_order_fill(
714 report: &FillReport,
715 state: &WsDispatchState,
716 emitter: &ExecutionEventEmitter,
717 ts_init: UnixNanos,
718) -> DispatchOutcome {
719 let Some(client_order_id) = report.client_order_id else {
720 return DispatchOutcome::External;
721 };
722
723 if state.filled_orders.contains(&client_order_id) {
724 log::debug!(
725 "Skipping stale fill for filled order: cid={client_order_id}, trade_id={}",
726 report.trade_id,
727 );
728 return DispatchOutcome::Skip;
729 }
730
731 let Some(mut context) = state.lookup_context(&client_order_id) else {
732 return DispatchOutcome::External;
733 };
734
735 let mut promoted_corrective: Option<(Quantity, HyperliquidExchangePlaceOrderRequest)> = None;
737
738 if state.has_pending_modify(&client_order_id)
741 && let Some(cached_voi) = state.cached_venue_order_id(&client_order_id)
742 && report.venue_order_id != cached_voi
743 {
744 let target = state.pending_modify_target_qty(&client_order_id);
745 let sent_request = state.modify_request(&client_order_id);
746 let price = sent_request
748 .as_ref()
749 .zip(context.price)
750 .and_then(|(r, cached)| Price::from_decimal_dp(r.price, cached.precision).ok())
751 .or(context.price);
752 let Some(price) = price else {
753 log::warn!(
754 "Cannot promote cancel-replace for {client_order_id} from fill: no target \
755 or cached price; buffering until the replacement ACCEPTED arrives",
756 );
757 state.buffer_fill(client_order_id, report.clone());
758 return DispatchOutcome::Tracked;
759 };
760 let updated_quantity = target.unwrap_or(context.quantity);
761 promote_cancel_replace(
762 client_order_id,
763 &context,
764 state,
765 emitter,
766 report.venue_order_id,
767 report.account_id,
768 price,
769 updated_quantity,
770 None,
771 report.ts_event,
772 ts_init,
773 );
774 if let Some(updated) = state.lookup_context(&client_order_id) {
776 context = updated;
777 }
778
779 if let (Some(target), Some(sent_request)) = (target, sent_request) {
780 promoted_corrective = Some((target, sent_request));
781 }
782 }
783
784 if state.check_and_insert_trade(report.trade_id) {
785 log::debug!(
786 "Skipping duplicate fill for {client_order_id}: trade_id={}",
787 report.trade_id
788 );
789 return DispatchOutcome::Tracked;
790 }
791
792 let previous = state
793 .previous_filled_qty(&client_order_id)
794 .unwrap_or_else(|| Quantity::zero(report.last_qty.precision));
795 let cumulative = previous + report.last_qty;
796
797 let is_terminal_fill = cumulative >= context.quantity;
798 if is_terminal_fill && !claim_terminal_order(client_order_id, state, OrderStatus::Filled) {
799 return DispatchOutcome::Skip;
800 }
801
802 ensure_accepted_emitted(
803 client_order_id,
804 report.venue_order_id,
805 report.account_id,
806 &context,
807 state,
808 emitter,
809 report.ts_event,
810 ts_init,
811 );
812
813 let filled = OrderFilled::new(
814 emitter.trader_id(),
815 context.identity.strategy_id,
816 context.identity.instrument_id,
817 client_order_id,
818 report.venue_order_id,
819 report.account_id,
820 report.trade_id,
821 context.identity.order_side,
822 context.identity.order_type,
823 report.last_qty,
824 report.last_px,
825 report.commission.currency,
826 report.liquidity_side,
827 UUID4::new(),
828 report.ts_event,
829 ts_init,
830 false,
831 report.venue_position_id,
832 Some(report.commission),
833 None,
834 );
835 emitter.send_order_event(OrderEventAny::Filled(filled));
836
837 state.record_filled_qty(client_order_id, cumulative);
838
839 if let Some((target, sent_request)) = promoted_corrective {
841 maybe_queue_corrective_reduce(
842 state,
843 client_order_id,
844 report.venue_order_id,
845 target,
846 sent_request,
847 );
848 }
849
850 if is_terminal_fill {
851 state.cleanup_terminal(&client_order_id);
852 }
853
854 DispatchOutcome::Tracked
855}
856
857fn handle_accepted(
858 report: &OrderStatusReport,
859 client_order_id: ClientOrderId,
860 context: &OrderContext,
861 state: &WsDispatchState,
862 emitter: &ExecutionEventEmitter,
863 ts_init: UnixNanos,
864) -> DispatchOutcome {
865 let venue_order_id = report.venue_order_id;
866 let ts_event = report.ts_last;
867 let account_id = report.account_id;
868
869 if let Some(cached_voi) = state.cached_venue_order_id(&client_order_id)
874 && cached_voi != venue_order_id
875 {
876 let price = report.price.or(context.price);
877 let Some(price) = price else {
878 log::warn!(
879 "Cannot emit OrderUpdated for cancel-replace {client_order_id}: \
880 no price on report and no cached price on context",
881 );
882 return DispatchOutcome::Skip;
883 };
884
885 let target_total_qty = state.pending_modify_target_qty(&client_order_id);
888 let updated_quantity = target_total_qty.unwrap_or(report.quantity);
889 let sent_request = state.modify_request(&client_order_id);
890
891 promote_cancel_replace(
892 client_order_id,
893 context,
894 state,
895 emitter,
896 venue_order_id,
897 account_id,
898 price,
899 updated_quantity,
900 report.trigger_price,
901 ts_event,
902 ts_init,
903 );
904
905 if let (Some(target), Some(sent_request)) = (target_total_qty, sent_request) {
906 maybe_queue_corrective_reduce(
907 state,
908 client_order_id,
909 venue_order_id,
910 target,
911 sent_request,
912 );
913 }
914
915 return DispatchOutcome::Tracked;
916 }
917
918 if state.emitted_accepted.contains(&client_order_id) {
919 state.update_context_price(&client_order_id, report.price);
923 return DispatchOutcome::Tracked;
924 }
925
926 state.insert_accepted(client_order_id);
927 state.record_venue_order_id(client_order_id, venue_order_id);
928 state.update_context_price(&client_order_id, report.price);
929
930 let accepted = OrderAccepted::new(
931 emitter.trader_id(),
932 context.identity.strategy_id,
933 context.identity.instrument_id,
934 client_order_id,
935 venue_order_id,
936 account_id,
937 UUID4::new(),
938 ts_event,
939 ts_init,
940 false,
941 );
942 emitter.send_order_event(OrderEventAny::Accepted(accepted));
943 DispatchOutcome::Tracked
944}
945
946#[allow(
949 clippy::too_many_arguments,
950 reason = "promotion needs the full OrderUpdated field set, sourced from two report shapes"
951)]
952fn promote_cancel_replace(
953 client_order_id: ClientOrderId,
954 context: &OrderContext,
955 state: &WsDispatchState,
956 emitter: &ExecutionEventEmitter,
957 venue_order_id: VenueOrderId,
958 account_id: AccountId,
959 price: Price,
960 quantity: Quantity,
961 trigger_price: Option<Price>,
962 ts_event: UnixNanos,
963 ts_init: UnixNanos,
964) {
965 state.record_venue_order_id(client_order_id, venue_order_id);
966 state.update_context_quantity(&client_order_id, quantity);
967 state.update_context_price(&client_order_id, Some(price));
968 state.claim_front_modify(&client_order_id, venue_order_id);
970
971 let updated = OrderUpdated::new(
972 emitter.trader_id(),
973 context.identity.strategy_id,
974 context.identity.instrument_id,
975 client_order_id,
976 quantity,
977 UUID4::new(),
978 ts_event,
979 ts_init,
980 false,
981 Some(venue_order_id),
982 Some(account_id),
983 Some(price),
984 trigger_price,
985 None,
986 false,
987 );
988 emitter.send_order_event(OrderEventAny::Updated(updated));
989
990 let buffered = state.drain_buffered_fills(&client_order_id);
993 for fill in buffered {
994 dispatch_order_fill(&fill, state, emitter, ts_init);
995 }
996}
997
998pub fn promote_replacement_from_query(
1005 report: &OrderStatusReport,
1006 state: &WsDispatchState,
1007 emitter: &ExecutionEventEmitter,
1008 ts_init: UnixNanos,
1009) -> bool {
1010 if report.order_status != OrderStatus::Accepted {
1011 return false;
1012 }
1013
1014 let Some(client_order_id) = report.client_order_id else {
1015 return false;
1016 };
1017
1018 if !state.has_pending_modify(&client_order_id) {
1019 return false;
1020 }
1021
1022 let Some(cached_voi) = state.cached_venue_order_id(&client_order_id) else {
1023 return false;
1024 };
1025
1026 if report.venue_order_id == cached_voi {
1027 return false;
1028 }
1029
1030 let Some(context) = state.lookup_context(&client_order_id) else {
1031 return false;
1032 };
1033
1034 let Some(price) = report.price.or(context.price) else {
1035 log::warn!(
1036 "Cannot promote cancel-replace from query for {client_order_id}: \
1037 no price on report and no cached price on context",
1038 );
1039 return false;
1040 };
1041
1042 let updated_quantity = state
1044 .pending_modify_target_qty(&client_order_id)
1045 .unwrap_or(report.quantity);
1046
1047 promote_cancel_replace(
1048 client_order_id,
1049 &context,
1050 state,
1051 emitter,
1052 report.venue_order_id,
1053 report.account_id,
1054 price,
1055 updated_quantity,
1056 report.trigger_price,
1057 report.ts_last,
1058 ts_init,
1059 );
1060
1061 log::debug!("Promoted cancel-replace replacement for {client_order_id} from query");
1062
1063 true
1064}
1065
1066fn maybe_queue_corrective_reduce(
1069 state: &WsDispatchState,
1070 client_order_id: ClientOrderId,
1071 venue_order_id: VenueOrderId,
1072 target: Quantity,
1073 sent_request: HyperliquidExchangePlaceOrderRequest,
1074) {
1075 let Ok(new_oid) = venue_order_id.as_str().parse::<u64>() else {
1076 return;
1077 };
1078
1079 let filled = state
1080 .previous_filled_qty(&client_order_id)
1081 .unwrap_or_else(|| Quantity::zero(target.precision));
1082 if filled >= target {
1083 return;
1084 }
1085
1086 let remaining = (target - filled).as_decimal().normalize();
1087
1088 let sent_size = sent_request.size;
1089 if sent_size > remaining {
1090 let mut corrective = sent_request;
1091 corrective.size = remaining;
1092
1093 state.mark_pending_modify(client_order_id, venue_order_id, target);
1094 state.stash_modify_request(client_order_id, corrective.clone());
1095 state.queue_corrective(client_order_id, new_oid, corrective);
1096
1097 log::warn!(
1098 "Cancel-replace left {client_order_id} oversized on {venue_order_id} \
1099 (sent {sent_size}, remaining {remaining}); queuing corrective reduce",
1100 );
1101 }
1102}
1103
1104fn handle_triggered(
1105 report: &OrderStatusReport,
1106 client_order_id: ClientOrderId,
1107 context: &OrderContext,
1108 state: &WsDispatchState,
1109 emitter: &ExecutionEventEmitter,
1110 ts_init: UnixNanos,
1111) -> DispatchOutcome {
1112 if !matches!(
1113 context.identity.order_type,
1114 OrderType::StopLimit | OrderType::TrailingStopLimit | OrderType::LimitIfTouched
1115 ) {
1116 log::debug!(
1117 "Ignoring TRIGGERED status for non-triggerable order type {:?}: {client_order_id}",
1118 context.identity.order_type,
1119 );
1120 return DispatchOutcome::Tracked;
1121 }
1122
1123 ensure_accepted_emitted(
1124 client_order_id,
1125 report.venue_order_id,
1126 report.account_id,
1127 context,
1128 state,
1129 emitter,
1130 report.ts_last,
1131 ts_init,
1132 );
1133
1134 let triggered = OrderTriggered::new(
1135 emitter.trader_id(),
1136 context.identity.strategy_id,
1137 context.identity.instrument_id,
1138 client_order_id,
1139 UUID4::new(),
1140 report.ts_last,
1141 ts_init,
1142 false,
1143 Some(report.venue_order_id),
1144 Some(report.account_id),
1145 );
1146 emitter.send_order_event(OrderEventAny::Triggered(triggered));
1147 DispatchOutcome::Tracked
1148}
1149
1150fn handle_canceled(
1151 report: &OrderStatusReport,
1152 client_order_id: ClientOrderId,
1153 context: &OrderContext,
1154 state: &WsDispatchState,
1155 emitter: &ExecutionEventEmitter,
1156 ts_init: UnixNanos,
1157) -> DispatchOutcome {
1158 let venue_order_id = report.venue_order_id;
1159
1160 if let Some(cached_voi) = state.cached_venue_order_id(&client_order_id)
1164 && cached_voi != venue_order_id
1165 {
1166 log::debug!(
1167 "Skipping stale CANCELED for {venue_order_id} (cached {cached_voi}) on {client_order_id}",
1168 );
1169 return DispatchOutcome::Skip;
1170 }
1171
1172 if state.pending_modify_contains_old(&client_order_id, venue_order_id) {
1178 log::debug!(
1179 "Skipping cancel-before-accept leg for {client_order_id}: venue_order_id={venue_order_id}",
1180 );
1181 return DispatchOutcome::Skip;
1182 }
1183
1184 if !claim_terminal_order(client_order_id, state, report.order_status) {
1185 return DispatchOutcome::Skip;
1186 }
1187
1188 ensure_accepted_emitted(
1189 client_order_id,
1190 venue_order_id,
1191 report.account_id,
1192 context,
1193 state,
1194 emitter,
1195 report.ts_last,
1196 ts_init,
1197 );
1198
1199 let canceled = OrderCanceled::new(
1200 emitter.trader_id(),
1201 context.identity.strategy_id,
1202 context.identity.instrument_id,
1203 client_order_id,
1204 UUID4::new(),
1205 report.ts_last,
1206 ts_init,
1207 false,
1208 Some(venue_order_id),
1209 Some(report.account_id),
1210 );
1211 emitter.send_order_event(OrderEventAny::Canceled(canceled));
1212
1213 state.cleanup_terminal(&client_order_id);
1214 DispatchOutcome::Tracked
1215}
1216
1217fn handle_expired(
1218 report: &OrderStatusReport,
1219 client_order_id: ClientOrderId,
1220 context: &OrderContext,
1221 state: &WsDispatchState,
1222 emitter: &ExecutionEventEmitter,
1223 ts_init: UnixNanos,
1224) -> DispatchOutcome {
1225 if !claim_terminal_order(client_order_id, state, report.order_status) {
1226 return DispatchOutcome::Skip;
1227 }
1228
1229 ensure_accepted_emitted(
1230 client_order_id,
1231 report.venue_order_id,
1232 report.account_id,
1233 context,
1234 state,
1235 emitter,
1236 report.ts_last,
1237 ts_init,
1238 );
1239
1240 let expired = OrderExpired::new(
1241 emitter.trader_id(),
1242 context.identity.strategy_id,
1243 context.identity.instrument_id,
1244 client_order_id,
1245 UUID4::new(),
1246 report.ts_last,
1247 ts_init,
1248 false,
1249 Some(report.venue_order_id),
1250 Some(report.account_id),
1251 );
1252 emitter.send_order_event(OrderEventAny::Expired(expired));
1253 state.cleanup_terminal(&client_order_id);
1254 DispatchOutcome::Tracked
1255}
1256
1257fn handle_rejected(
1258 report: &OrderStatusReport,
1259 client_order_id: ClientOrderId,
1260 context: &OrderContext,
1261 state: &WsDispatchState,
1262 emitter: &ExecutionEventEmitter,
1263 ts_init: UnixNanos,
1264) -> DispatchOutcome {
1265 if state.submission_pending(&client_order_id) {
1266 state.buffer_submission_rejection(client_order_id, report.clone());
1267 return DispatchOutcome::Skip;
1268 }
1269
1270 if !claim_terminal_order(client_order_id, state, report.order_status) {
1271 return DispatchOutcome::Skip;
1272 }
1273
1274 let reason = report
1275 .cancel_reason
1276 .clone()
1277 .unwrap_or_else(|| "Order rejected by exchange".to_string());
1278 let rejected = OrderRejected::new(
1279 emitter.trader_id(),
1280 context.identity.strategy_id,
1281 context.identity.instrument_id,
1282 client_order_id,
1283 report.account_id,
1284 Ustr::from(&reason),
1285 UUID4::new(),
1286 report.ts_last,
1287 ts_init,
1288 false,
1289 report.post_only && reason.contains(HYPERLIQUID_POST_ONLY_WOULD_MATCH),
1290 );
1291 emitter.send_order_event(OrderEventAny::Rejected(rejected));
1292 state.cleanup_terminal(&client_order_id);
1293 DispatchOutcome::Tracked
1294}
1295
1296fn claim_terminal_order(
1297 client_order_id: ClientOrderId,
1298 state: &WsDispatchState,
1299 status: OrderStatus,
1300) -> bool {
1301 let claimed = state.insert_filled(client_order_id);
1302 if !claimed {
1303 log::debug!("Skipping duplicate terminal event for {client_order_id}: status={status:?}",);
1304 }
1305
1306 claimed
1307}
1308
1309fn handle_filled_marker(
1310 _client_order_id: ClientOrderId,
1311 _state: &WsDispatchState,
1312) -> DispatchOutcome {
1313 DispatchOutcome::Tracked
1321}
1322
1323#[allow(clippy::too_many_arguments)]
1331fn ensure_accepted_emitted(
1332 client_order_id: ClientOrderId,
1333 venue_order_id: VenueOrderId,
1334 account_id: AccountId,
1335 context: &OrderContext,
1336 state: &WsDispatchState,
1337 emitter: &ExecutionEventEmitter,
1338 ts_event: UnixNanos,
1339 ts_init: UnixNanos,
1340) {
1341 if state.emitted_accepted.contains(&client_order_id) {
1342 return;
1343 }
1344 state.insert_accepted(client_order_id);
1345 state.record_venue_order_id(client_order_id, venue_order_id);
1346
1347 let accepted = OrderAccepted::new(
1348 emitter.trader_id(),
1349 context.identity.strategy_id,
1350 context.identity.instrument_id,
1351 client_order_id,
1352 venue_order_id,
1353 account_id,
1354 UUID4::new(),
1355 ts_event,
1356 ts_init,
1357 false,
1358 );
1359 emitter.send_order_event(OrderEventAny::Accepted(accepted));
1360}
1361
1362#[cfg(test)]
1363mod tests {
1364 use nautilus_live::execution::context::OrderIdentity;
1365 use nautilus_model::{
1366 enums::{OrderSide, TimeInForce},
1367 identifiers::{ClientOrderId, InstrumentId, StrategyId, TradeId},
1368 };
1369 use rstest::rstest;
1370 use rust_decimal::Decimal;
1371
1372 use super::*;
1373 use crate::http::models::{
1374 HyperliquidExchangeLimitParams, HyperliquidExchangeOrderKind, HyperliquidExchangeTif,
1375 };
1376
1377 fn make_context(client_order_id: ClientOrderId) -> OrderContext {
1378 OrderContext {
1379 identity: OrderIdentity {
1380 client_order_id,
1381 strategy_id: StrategyId::from("S-001"),
1382 instrument_id: InstrumentId::from("BTC-USD-PERP.HYPERLIQUID"),
1383 order_side: OrderSide::Buy,
1384 order_type: OrderType::Limit,
1385 },
1386 quantity: Quantity::from("0.0001"),
1387 price: None,
1388 trigger_price: None,
1389 trigger_type: None,
1390 time_in_force: TimeInForce::Gtc,
1391 is_post_only: false,
1392 is_reduce_only: false,
1393 is_quote_quantity: false,
1394 }
1395 }
1396
1397 #[rstest]
1398 fn test_register_and_lookup_context() {
1399 let state = WsDispatchState::new();
1400 let cid = ClientOrderId::new("O-001");
1401 state.register_context(make_context(cid));
1402
1403 assert_eq!(state.lookup_context(&cid), Some(make_context(cid)));
1404 }
1405
1406 #[rstest]
1407 fn test_lookup_context_missing_returns_none() {
1408 let state = WsDispatchState::new();
1409 let cid = ClientOrderId::new("not-tracked");
1410 assert!(state.lookup_context(&cid).is_none());
1411 }
1412
1413 #[rstest]
1414 fn test_update_context_refreshes_price_and_quantity() {
1415 let state = WsDispatchState::new();
1416 let cid = ClientOrderId::new("O-003");
1417 state.register_context(make_context(cid));
1418
1419 state.update_context_price(&cid, Some(Price::from("56731.5")));
1420 state.update_context_quantity(&cid, Quantity::from("0.0002"));
1421
1422 assert_eq!(
1423 state.lookup_context(&cid),
1424 Some(OrderContext {
1425 price: Some(Price::from("56731.5")),
1426 quantity: Quantity::from("0.0002"),
1427 ..make_context(cid)
1428 }),
1429 );
1430 }
1431
1432 #[rstest]
1433 fn test_update_context_price_without_price_keeps_cached_price() {
1434 let state = WsDispatchState::new();
1435 let cid = ClientOrderId::new("O-004");
1436 let cached = OrderContext {
1437 price: Some(Price::from("56730.0")),
1438 ..make_context(cid)
1439 };
1440 state.register_context(cached);
1441
1442 state.update_context_price(&cid, None);
1443
1444 assert_eq!(state.lookup_context(&cid), Some(cached));
1445 }
1446
1447 #[rstest]
1448 fn test_insert_accepted_dedup() {
1449 let state = WsDispatchState::new();
1450 let cid = ClientOrderId::new("O-002");
1451 assert!(!state.emitted_accepted.contains(&cid));
1452 state.insert_accepted(cid);
1453 assert!(state.emitted_accepted.contains(&cid));
1454 state.insert_accepted(cid);
1455 assert!(state.emitted_accepted.contains(&cid));
1456 }
1457
1458 #[rstest]
1459 fn test_check_and_insert_trade_detects_duplicates() {
1460 let state = WsDispatchState::new();
1461 let trade = TradeId::new("trade-1");
1462 assert!(!state.check_and_insert_trade(trade));
1463 assert!(state.check_and_insert_trade(trade));
1464 }
1465
1466 #[rstest]
1467 fn test_pending_modify_roundtrip() {
1468 let state = WsDispatchState::new();
1469 let cid = ClientOrderId::new("O-010");
1470 let voi = VenueOrderId::new("v-1");
1471 let target_qty = Quantity::from("0.0001");
1472
1473 assert!(state.pending_modify(&cid).is_none());
1474 assert!(state.pending_modify_target_qty(&cid).is_none());
1475 state.mark_pending_modify(cid, voi, target_qty);
1476 assert_eq!(state.pending_modify(&cid), Some(voi));
1477 assert_eq!(state.pending_modify_target_qty(&cid), Some(target_qty));
1478 state.clear_pending_modify(&cid);
1479 assert!(state.pending_modify(&cid).is_none());
1480 assert!(state.pending_modify_target_qty(&cid).is_none());
1481 }
1482
1483 #[rstest]
1484 fn test_cleanup_terminal_preserves_filled_marker() {
1485 let state = WsDispatchState::new();
1486 let cid = ClientOrderId::new("O-020");
1487 state.register_context(make_context(cid));
1488 state.insert_accepted(cid);
1489 state.mark_pending_modify(cid, VenueOrderId::new("v-1"), Quantity::from("0.0001"));
1490 state.insert_filled(cid);
1491 state.cleanup_terminal(&cid);
1492
1493 assert!(state.lookup_context(&cid).is_none());
1494 assert!(!state.emitted_accepted.contains(&cid));
1495 assert!(state.pending_modify(&cid).is_none());
1496 assert!(state.pending_modify_target_qty(&cid).is_none());
1497 assert!(state.filled_orders.contains(&cid));
1499 }
1500
1501 #[rstest]
1502 fn test_cleanup_terminal_clears_corrective_state() {
1503 let state = WsDispatchState::new();
1504 let cid = ClientOrderId::new("O-021");
1505 let request = sample_request(Decimal::from(1));
1506 state.mark_pending_modify(cid, VenueOrderId::new("v-1"), Quantity::from("1"));
1507 state.stash_modify_request(cid, request.clone());
1508 state.queue_corrective(cid, 1, request);
1509 assert!(state.modify_request(&cid).is_some());
1510
1511 state.cleanup_terminal(&cid);
1512
1513 assert!(state.modify_request(&cid).is_none());
1514 assert!(state.take_corrective(&cid).is_none());
1515 assert!(state.pending_modify(&cid).is_none());
1516 }
1517
1518 fn sample_request(size: Decimal) -> HyperliquidExchangePlaceOrderRequest {
1519 HyperliquidExchangePlaceOrderRequest {
1520 asset: 0,
1521 is_buy: true,
1522 price: "100".parse::<Decimal>().unwrap(),
1523 size,
1524 reduce_only: false,
1525 kind: HyperliquidExchangeOrderKind::Limit {
1526 limit: HyperliquidExchangeLimitParams {
1527 tif: HyperliquidExchangeTif::Gtc,
1528 },
1529 },
1530 cloid: None,
1531 }
1532 }
1533
1534 #[rstest]
1535 fn test_modify_chain_keeps_both_intents_on_rapid_modifies() {
1536 let state = WsDispatchState::new();
1537 let cid = ClientOrderId::new("O-100");
1538 let g0 =
1539 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00020"));
1540 let g1 =
1541 state.mark_pending_modify(cid, VenueOrderId::new("v-1"), Quantity::from("0.00030"));
1542
1543 assert_ne!(g0, g1);
1544 assert!(state.has_pending_modify(&cid));
1545 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new("v-0")));
1547 assert_eq!(
1548 state.pending_modify_target_qty(&cid),
1549 Some(Quantity::from("0.00020")),
1550 );
1551 assert!(state.pending_modify_contains_old(&cid, VenueOrderId::new("v-0")));
1553 assert!(state.pending_modify_contains_old(&cid, VenueOrderId::new("v-1")));
1554 }
1555
1556 #[rstest]
1557 fn test_clear_modify_generation_preserves_newer_intent() {
1558 let state = WsDispatchState::new();
1559 let cid = ClientOrderId::new("O-101");
1560 let g0 =
1563 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00020"));
1564 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00030"));
1565
1566 state.clear_modify_generation(&cid, g0);
1569
1570 assert!(state.has_pending_modify(&cid));
1571 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new("v-0")));
1572 assert_eq!(
1573 state.pending_modify_target_qty(&cid),
1574 Some(Quantity::from("0.00030")),
1575 );
1576 assert!(state.pending_modify_contains_old(&cid, VenueOrderId::new("v-0")));
1577 }
1578
1579 #[rstest]
1580 fn test_claim_front_modify_advances_next_old_id() {
1581 let state = WsDispatchState::new();
1582 let cid = ClientOrderId::new("O-102");
1583 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00020"));
1585 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00030"));
1586
1587 let claimed = state.claim_front_modify(&cid, VenueOrderId::new("v-1"));
1590 assert_eq!(
1591 claimed.map(|i| i.target_qty),
1592 Some(Quantity::from("0.00020"))
1593 );
1594
1595 assert!(state.has_pending_modify(&cid));
1596 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new("v-1")));
1597 assert!(state.pending_modify_contains_old(&cid, VenueOrderId::new("v-1")));
1598 assert!(!state.pending_modify_contains_old(&cid, VenueOrderId::new("v-0")));
1600
1601 let claimed2 = state.claim_front_modify(&cid, VenueOrderId::new("v-2"));
1603 assert_eq!(
1604 claimed2.map(|i| i.target_qty),
1605 Some(Quantity::from("0.00030"))
1606 );
1607 assert!(!state.has_pending_modify(&cid));
1608 assert!(state.pending_modify(&cid).is_none());
1609 }
1610
1611 #[rstest]
1612 fn test_modify_chain_caps_and_evicts_oldest() {
1613 let state = WsDispatchState::new();
1614 let cid = ClientOrderId::new("O-106");
1615 for i in 0..=MAX_PENDING_MODIFY_INTENTS {
1617 let voi = format!("v-{i}");
1618 state.mark_pending_modify(cid, VenueOrderId::new(&voi), Quantity::from("0.00020"));
1619 }
1620
1621 assert!(!state.pending_modify_contains_old(&cid, VenueOrderId::new("v-0")));
1624 let newest = format!("v-{MAX_PENDING_MODIFY_INTENTS}");
1625 assert!(state.pending_modify_contains_old(&cid, VenueOrderId::new(&newest)));
1626 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new("v-1")));
1627 }
1628
1629 #[rstest]
1630 fn test_clear_front_modify_reparents_next_old() {
1631 let state = WsDispatchState::new();
1632 let cid = ClientOrderId::new("O-104");
1633 let g_front =
1636 state.mark_pending_modify(cid, VenueOrderId::new("v-1"), Quantity::from("0.00020"));
1637 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00030"));
1638
1639 state.clear_modify_generation(&cid, g_front);
1643
1644 assert!(state.has_pending_modify(&cid));
1645 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new("v-1")));
1646 assert!(state.pending_modify_contains_old(&cid, VenueOrderId::new("v-1")));
1647 assert!(!state.pending_modify_contains_old(&cid, VenueOrderId::new("v-0")));
1648 }
1649
1650 #[rstest]
1651 fn test_clear_non_front_modify_leaves_front_old() {
1652 let state = WsDispatchState::new();
1653 let cid = ClientOrderId::new("O-105");
1654 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00020"));
1655 let g_back =
1656 state.mark_pending_modify(cid, VenueOrderId::new("v-9"), Quantity::from("0.00030"));
1657
1658 state.clear_modify_generation(&cid, g_back);
1660
1661 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new("v-0")));
1662 assert!(!state.pending_modify_contains_old(&cid, VenueOrderId::new("v-9")));
1663 }
1664
1665 #[rstest]
1666 fn test_stash_modify_request_targets_latest_intent() {
1667 let state = WsDispatchState::new();
1668 let cid = ClientOrderId::new("O-103");
1669 state.mark_pending_modify(cid, VenueOrderId::new("v-0"), Quantity::from("0.00020"));
1670 state.stash_modify_request(cid, sample_request(Decimal::from(1)));
1671 state.mark_pending_modify(cid, VenueOrderId::new("v-1"), Quantity::from("0.00030"));
1672 state.stash_modify_request(cid, sample_request(Decimal::from(2)));
1673
1674 assert_eq!(
1676 state.modify_request(&cid).map(|r| r.size),
1677 Some(Decimal::from(1)),
1678 );
1679 state.claim_front_modify(&cid, VenueOrderId::new("v-1"));
1681 assert_eq!(
1682 state.modify_request(&cid).map(|r| r.size),
1683 Some(Decimal::from(2)),
1684 );
1685 }
1686}