1use std::{
19 sync::{Arc, Mutex},
20 time::{Duration, Instant},
21};
22
23use ahash::AHashMap;
24use anyhow::Context;
25use async_trait::async_trait;
26use nautilus_common::{
27 cache::fifo::FifoCache,
28 clients::{ExecutionClient, SocketReconnectRegistry},
29 live::{runner::get_exec_event_sender, runtime::get_runtime, task::TaskHandles},
30 messages::execution::{
31 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
32 GenerateOrderStatusReport, GenerateOrderStatusReports, GeneratePositionStatusReports,
33 ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
34 },
35};
36use nautilus_core::{
37 MUTEX_POISONED, Params, UUID4, UnixNanos,
38 time::{AtomicTime, get_atomic_clock_realtime},
39};
40use nautilus_live::{ExecutionClientCore, ExecutionEventEmitter, execution::context::OrderContext};
41use nautilus_model::{
42 accounts::AccountAny,
43 enums::{AccountType, OmsType, OrderSide, OrderStatus, OrderType},
44 identifiers::{
45 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, Venue, VenueOrderId,
46 },
47 orders::{Order, any::OrderAny},
48 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
49 types::{AccountBalance, MarginBalance, Quantity},
50};
51
52#[derive(Debug, Clone)]
53struct StagedBracketChild {
54 order: OrderAny,
55 request: HyperliquidExecPlaceOrderRequest,
56}
57
58#[derive(Debug, Default)]
59struct StagedBracketState {
60 children_by_parent: AHashMap<ClientOrderId, Vec<StagedBracketChild>>,
61 active_children: AHashMap<ClientOrderId, StagedBracketChild>,
62 active_siblings: AHashMap<ClientOrderId, ClientOrderId>,
63}
64
65impl StagedBracketState {
66 fn stage(&mut self, parent_id: ClientOrderId, children: Vec<StagedBracketChild>) {
67 self.children_by_parent.insert(parent_id, children);
68 }
69
70 fn activate(&mut self, parent_id: &ClientOrderId) -> Option<Vec<StagedBracketChild>> {
71 let children = self.children_by_parent.remove(parent_id)?;
72 self.track_active(&children);
73
74 Some(children)
75 }
76
77 fn restore_active(&mut self, children: &[StagedBracketChild]) {
78 self.track_active(children);
79 }
80
81 fn track_active(&mut self, children: &[StagedBracketChild]) {
82 let child_ids = children
83 .iter()
84 .map(|child| child.order.client_order_id())
85 .collect::<Vec<_>>();
86
87 for child in children {
88 let child_id = child.order.client_order_id();
89 if let Some(sibling_id) = child
90 .order
91 .linked_order_ids()
92 .and_then(|ids| ids.iter().find(|id| child_ids.contains(id)))
93 {
94 self.active_siblings.insert(child_id, *sibling_id);
95 }
96 self.active_children.insert(child_id, child.clone());
97 }
98 }
99
100 fn contains_parent(&self, parent_id: &ClientOrderId) -> bool {
101 self.children_by_parent.contains_key(parent_id)
102 }
103
104 fn cancel_child(&mut self, child_id: &ClientOrderId) -> Option<OrderAny> {
105 let parent_id = self
106 .children_by_parent
107 .iter()
108 .find_map(|(parent_id, children)| {
109 children
110 .iter()
111 .any(|child| child.order.client_order_id() == *child_id)
112 .then_some(*parent_id)
113 })?;
114 let children = self.children_by_parent.get_mut(&parent_id)?;
115 let index = children
116 .iter()
117 .position(|child| child.order.client_order_id() == *child_id)?;
118 let child = children.remove(index);
119
120 if children.is_empty() {
121 self.children_by_parent.remove(&parent_id);
122 }
123
124 Some(child.order)
125 }
126
127 fn cancel_for_parent(&mut self, parent_id: &ClientOrderId) -> Vec<OrderAny> {
128 self.children_by_parent
129 .remove(parent_id)
130 .map(|children| children.into_iter().map(|child| child.order).collect())
131 .unwrap_or_default()
132 }
133
134 fn take_active_sibling(
135 &mut self,
136 client_order_id: &ClientOrderId,
137 ) -> Option<StagedBracketChild> {
138 self.active_children.remove(client_order_id);
139 let sibling_id = self.active_siblings.remove(client_order_id)?;
140 self.active_siblings.remove(&sibling_id);
141 self.active_children.remove(&sibling_id)
142 }
143
144 fn active_sibling(&self, client_order_id: &ClientOrderId) -> Option<StagedBracketChild> {
145 self.active_siblings
146 .get(client_order_id)
147 .and_then(|sibling_id| self.active_children.get(sibling_id))
148 .cloned()
149 }
150}
151use tokio::task::JoinHandle;
152use ustr::Ustr;
153
154use crate::{
155 account::resolve_execution_account_address,
156 common::{
157 consts::{
158 HYPERLIQUID_BUILDER_APPROVAL_DOCS_URL, HYPERLIQUID_BUILDER_FEE_NOT_APPROVED,
159 HYPERLIQUID_POST_ONLY_WOULD_MATCH, HYPERLIQUID_VENUE,
160 },
161 credential::Secrets,
162 enums::HyperliquidProductType,
163 parse::{
164 clamp_price_to_precision, derive_limit_from_trigger, derive_market_order_price,
165 extract_error_message, extract_inner_error, extract_inner_errors, normalize_price,
166 order_to_hyperliquid_request_with_asset_and_cloid,
167 parse_combined_account_balances_and_margins, round_to_sig_figs,
168 },
169 socket::{SocketStatePublisher, USER_STREAMS_ENDPOINT},
170 },
171 config::HyperliquidExecClientConfig,
172 http::{
173 client::HyperliquidHttpClient,
174 models::{
175 ClearinghouseState, Cloid, HyperliquidExecAction, HyperliquidExecCancelByCloidRequest,
176 HyperliquidExecCancelOrderRequest, HyperliquidExecGrouping,
177 HyperliquidExecModifyOrderRequest, HyperliquidExecModifyTarget,
178 HyperliquidExecOrderKind, HyperliquidExecPlaceOrderRequest, HyperliquidExecTpSl,
179 SpotClearinghouseState,
180 },
181 parse::derive_outcome_settlements,
182 },
183 outcome_settlement::{OutcomeSettlementTracker, build_settlement_fills},
184 websocket::{
185 ExecutionReport, NautilusWsMessage,
186 client::HyperliquidWebSocketClient,
187 dispatch::{
188 DispatchOutcome, WsDispatchState, dispatch_order_event, dispatch_order_fill,
189 promote_replacement_from_query,
190 },
191 },
192};
193
194#[derive(Debug)]
195pub struct HyperliquidExecutionClient {
196 core: ExecutionClientCore,
197 clock: &'static AtomicTime,
198 config: HyperliquidExecClientConfig,
199 emitter: ExecutionEventEmitter,
200 http_client: HyperliquidHttpClient,
201 ws_client: HyperliquidWebSocketClient,
202 pending_tasks: TaskHandles,
203 ws_stream_handle: Option<JoinHandle<()>>,
204 settlement_poll_handle: Option<JoinHandle<()>>,
205 ws_dispatch_state: Arc<WsDispatchState>,
206 staged_brackets: Arc<Mutex<StagedBracketState>>,
207 outcome_settlement_tracker: Arc<Mutex<OutcomeSettlementTracker>>,
208 socket_registry: SocketReconnectRegistry,
209}
210
211impl HyperliquidExecutionClient {
212 pub fn config(&self) -> &HyperliquidExecClientConfig {
214 &self.config
215 }
216
217 #[must_use]
224 pub fn ws_dispatch_state(&self) -> &Arc<WsDispatchState> {
225 &self.ws_dispatch_state
226 }
227
228 #[allow(
236 clippy::missing_panics_doc,
237 reason = "pending_tasks mutex poisoning is not expected"
238 )]
239 #[must_use]
240 pub fn pending_tasks_all_finished(&self) -> bool {
241 self.pending_tasks.all_finished()
242 }
243
244 fn resolve_slippage_bps(&self, params: Option<&Params>) -> u32 {
245 params
246 .and_then(|p| p.get_u64("market_order_slippage_bps"))
247 .map_or(self.config.market_order_slippage_bps, |v| v as u32)
248 }
249
250 fn validate_order_submission(&self, order: &OrderAny) -> anyhow::Result<()> {
251 validate_order_for_hyperliquid(order)
252 }
253
254 fn order_request(
255 &self,
256 order: &OrderAny,
257 slippage_bps: u32,
258 ) -> anyhow::Result<HyperliquidExecPlaceOrderRequest> {
259 validate_order_for_hyperliquid(order)?;
260
261 let symbol = order.instrument_id().symbol.inner();
262 let asset = self
263 .http_client
264 .get_asset_index_for_symbol(symbol)
265 .with_context(|| format!("Asset index not found for {symbol}"))?;
266 let price_decimals = self
267 .http_client
268 .get_price_precision_for_symbol(symbol)
269 .unwrap_or(2);
270 let cloid = self
271 .http_client
272 .get_or_generate_client_order_id_cloid(order.client_order_id());
273 let mut request = order_to_hyperliquid_request_with_asset_and_cloid(
274 order,
275 asset,
276 price_decimals,
277 self.config.normalize_prices,
278 slippage_bps,
279 None,
280 )?;
281 request.cloid = Some(cloid);
282
283 Ok(request)
284 }
285
286 fn restore_staged_brackets(&self) -> Vec<ClientOrderId> {
287 let order_lists = self
288 .core
289 .cache()
290 .order_lists(Some(&self.core.venue), None, None, None)
291 .into_iter()
292 .cloned()
293 .collect::<Vec<_>>();
294 let mut ready_parent_ids = Vec::new();
295
296 for order_list in order_lists {
297 let orders = {
298 let cache = self.core.cache();
299 order_list
300 .client_order_ids
301 .iter()
302 .filter_map(|client_order_id| {
303 cache.order(client_order_id).map(|order| order.clone())
304 })
305 .collect::<Vec<_>>()
306 };
307
308 if orders.len() != order_list.client_order_ids.len()
309 || determine_order_list_grouping(&orders) != HyperliquidExecGrouping::NormalTpsl
310 {
311 continue;
312 }
313
314 let (mut orders, mut requests) = match orders
315 .iter()
316 .map(|order| self.order_request(order, self.config.market_order_slippage_bps))
317 .collect::<anyhow::Result<Vec<_>>>()
318 {
319 Ok(requests) => order_normal_tpsl_submission(
320 orders,
321 requests,
322 HyperliquidExecGrouping::NormalTpsl,
323 ),
324 Err(e) => {
325 log::warn!("Cannot restore staged bracket {}: {e}", order_list.id,);
326 continue;
327 }
328 };
329 let parent = orders.remove(0);
330 let parent_request = requests.remove(0);
331 let parent_id = parent.client_order_id();
332 let (staged_children, active_children): (Vec<_>, Vec<_>) = orders
333 .drain(..)
334 .zip(requests.drain(..))
335 .filter(|(order, _)| order.is_active_local())
336 .map(|(order, request)| StagedBracketChild { order, request })
337 .partition(|child| child.order.status() == OrderStatus::Initialized);
338
339 if (staged_children.is_empty() && active_children.is_empty())
340 || (!parent.is_open() && parent.filled_qty().raw == 0)
341 || self
342 .staged_brackets
343 .lock()
344 .expect(MUTEX_POISONED)
345 .contains_parent(&parent_id)
346 {
347 continue;
348 }
349
350 self.restore_order_context(&parent, &parent_request);
351 for child in &active_children {
352 self.restore_order_context(&child.order, &child.request);
353 }
354
355 let has_staged_children = !staged_children.is_empty();
356 let mut state = self.staged_brackets.lock().expect(MUTEX_POISONED);
357 if has_staged_children {
358 state.stage(parent_id, staged_children);
359 }
360 state.restore_active(&active_children);
361 drop(state);
362
363 if has_staged_children && parent.filled_qty().raw > 0 {
364 ready_parent_ids.push(parent_id);
365 }
366 }
367
368 if !ready_parent_ids.is_empty() {
369 log::info!(
370 "Restored {} staged bracket parent(s) with prior fills",
371 ready_parent_ids.len(),
372 );
373 }
374
375 ready_parent_ids
376 }
377
378 fn restore_order_context(&self, order: &OrderAny, request: &HyperliquidExecPlaceOrderRequest) {
379 let client_order_id = order.client_order_id();
380 let cloid = request.cloid.expect("order conversion must set a CLOID");
381 self.http_client
382 .cache_client_order_id_cloid(client_order_id, cloid);
383 self.ws_client
384 .cache_cloid_mapping(Ustr::from(&cloid.to_hex()), client_order_id);
385 self.ws_dispatch_state
386 .register_context(OrderContext::from(order));
387
388 if let Some(venue_order_id) = order.venue_order_id() {
389 self.ws_dispatch_state
390 .record_venue_order_id(client_order_id, venue_order_id);
391 self.ws_dispatch_state.insert_accepted(client_order_id);
392 }
393 }
394
395 pub fn new(
401 core: ExecutionClientCore,
402 config: HyperliquidExecClientConfig,
403 ) -> anyhow::Result<Self> {
404 let secrets = Secrets::resolve(
405 config.private_key.as_deref(),
406 config.vault_address.as_deref(),
407 config.environment,
408 )
409 .context("Hyperliquid execution client requires private key")?;
410
411 let account_address = resolve_execution_account_address(
412 config.private_key.as_deref(),
413 config.vault_address.as_deref(),
414 config.account_address.as_deref(),
415 config.environment,
416 )?;
417
418 let mut http_client = HyperliquidHttpClient::with_secrets(
419 &secrets,
420 config.http_timeout_secs,
421 config.proxy_url.clone(),
422 )
423 .context("failed to create Hyperliquid HTTP client")?;
424
425 http_client.set_account_id(core.account_id);
426 http_client.set_account_address(account_address);
427 http_client.set_normalize_prices(config.normalize_prices);
428 http_client.set_market_order_slippage_bps(config.market_order_slippage_bps);
429 http_client.set_include_builder_attribution(config.include_builder_attribution);
430
431 if let Some(url) = &config.base_url_http {
433 http_client.set_base_info_url(url.clone());
434 }
435
436 if let Some(url) = &config.base_url_exchange {
437 http_client.set_base_exchange_url(url.clone());
438 }
439
440 let ws_url = config.base_url_ws.clone();
441 let mut ws_client = HyperliquidWebSocketClient::new(
442 ws_url,
443 config.environment,
444 Some(core.account_id),
445 config.transport_backend,
446 config.proxy_url.clone(),
447 );
448 let socket_registry = SocketReconnectRegistry::default();
449 if let Some(publisher) = SocketStatePublisher::new(core.client_id, socket_registry.clone())
450 {
451 ws_client = ws_client.with_socket_control(publisher.control(USER_STREAMS_ENDPOINT));
452 }
453 ws_client.set_post_timeout(Duration::from_secs(config.ws_post_timeout_secs));
454
455 let clock = get_atomic_clock_realtime();
456 let emitter = ExecutionEventEmitter::new(
457 clock,
458 core.trader_id,
459 core.account_id,
460 AccountType::Margin,
461 None,
462 );
463
464 Ok(Self {
465 core,
466 clock,
467 config,
468 emitter,
469 http_client,
470 ws_client,
471 pending_tasks: TaskHandles::default(),
472 ws_stream_handle: None,
473 settlement_poll_handle: None,
474 ws_dispatch_state: Arc::new(WsDispatchState::new()),
475 staged_brackets: Arc::new(Mutex::new(StagedBracketState::default())),
476 outcome_settlement_tracker: Arc::new(Mutex::new(OutcomeSettlementTracker::new())),
477 socket_registry,
478 })
479 }
480
481 fn register_order_context(&self, order: &OrderAny) {
482 register_order_context_into(&self.ws_dispatch_state, order);
483 }
484
485 async fn ensure_instruments_initialized_async(&self) -> anyhow::Result<()> {
486 if self.core.instruments_initialized() {
487 return Ok(());
488 }
489
490 let instruments = self
491 .http_client
492 .request_instruments()
493 .await
494 .context("failed to request Hyperliquid instruments")?;
495
496 if instruments.is_empty() {
497 log::warn!(
498 "Instrument bootstrap yielded no instruments; WebSocket submissions may fail"
499 );
500 } else {
501 log::debug!("Initialized {} instruments", instruments.len());
502
503 for instrument in &instruments {
504 self.http_client.cache_instrument(instrument);
505 }
506 }
507
508 self.core.set_instruments_initialized();
509 Ok(())
510 }
511
512 async fn refresh_account_state(&self) -> anyhow::Result<()> {
513 let account_address = self.get_account_address()?;
514
515 let (perp_state, spot_state) = self
516 .fetch_combined_clearinghouse_state(&account_address)
517 .await?;
518
519 log::debug!(
520 "Received clearinghouse state: cross_margin_summary={:?}, asset_positions={}, spot_balances={}",
521 perp_state.cross_margin_summary,
522 perp_state.asset_positions.len(),
523 spot_state.balances.len(),
524 );
525
526 let (balances, margins) =
527 parse_combined_account_balances_and_margins(&perp_state, &spot_state)
528 .context("failed to parse combined account balances and margins")?;
529
530 let ts_event = self.clock.get_time_ns();
533 self.emitter
534 .emit_account_state(balances, margins, true, ts_event, None);
535
536 log::debug!("Account state updated successfully");
537 Ok(())
538 }
539
540 async fn fetch_combined_clearinghouse_state(
541 &self,
542 account_address: &str,
543 ) -> anyhow::Result<(ClearinghouseState, SpotClearinghouseState)> {
544 let perp_json = self
545 .http_client
546 .info_clearinghouse_state(account_address)
547 .await
548 .context("failed to fetch clearinghouse state")?;
549 let perp_state: ClearinghouseState = serde_json::from_value(perp_json)
550 .context("failed to deserialize clearinghouse state")?;
551
552 let spot_json = self
553 .http_client
554 .info_spot_clearinghouse_state(account_address)
555 .await
556 .context("failed to fetch spot clearinghouse state")?;
557 let spot_state: SpotClearinghouseState = serde_json::from_value(spot_json)
558 .context("failed to deserialize spot clearinghouse state")?;
559
560 Ok((perp_state, spot_state))
561 }
562
563 async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
564 let account_id = self.core.account_id;
565
566 if self.core.cache().account(&account_id).is_some() {
567 log::info!("Account {account_id} registered");
568 return Ok(());
569 }
570
571 let start = Instant::now();
572 let timeout = Duration::from_secs_f64(timeout_secs);
573 let interval = Duration::from_millis(10);
574
575 loop {
576 tokio::time::sleep(interval).await;
577
578 if self.core.cache().account(&account_id).is_some() {
579 log::info!("Account {account_id} registered");
580 return Ok(());
581 }
582
583 if start.elapsed() >= timeout {
584 anyhow::bail!(
585 "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
586 );
587 }
588 }
589 }
590
591 fn get_account_address(&self) -> anyhow::Result<String> {
592 self.http_client
593 .get_account_address()
594 .context("failed to get account address from HTTP client")
595 }
596
597 fn spawn_task<F>(&self, description: &'static str, fut: F)
598 where
599 F: std::future::Future<Output = anyhow::Result<()>> + Send + 'static,
600 {
601 let runtime = get_runtime();
602 let handle = runtime.spawn(async move {
603 if let Err(e) = fut.await {
604 log::warn!("{description} failed: {e:?}");
605 }
606 });
607
608 self.pending_tasks.push(handle);
609 }
610
611 fn start_outcome_settlement_poll(&mut self) -> anyhow::Result<()> {
612 let poll_secs = self.config.outcome_settlement_poll_secs;
613 if poll_secs == 0 {
614 log::debug!("Outcome settlement polling disabled by config");
615 return Ok(());
616 }
617
618 let http_client = self.http_client.clone();
619 let emitter = self.emitter.clone();
620 let tracker = self.outcome_settlement_tracker.clone();
621 let account_id = self.core.account_id;
622 let account_address = self.get_account_address()?;
623 let clock = self.clock;
624
625 let handle = get_runtime().spawn(async move {
628 let mut interval = tokio::time::interval(Duration::from_secs(poll_secs));
629 interval.tick().await;
630
631 loop {
632 interval.tick().await;
633
634 let meta = match http_client.get_outcome_meta().await {
635 Ok(meta) => meta,
636 Err(e) => {
637 log::warn!("Outcome meta poll failed: {e}");
638 continue;
639 }
640 };
641
642 let settlements = derive_outcome_settlements(&meta);
643 if settlements.is_empty() {
644 continue;
645 }
646
647 let spot_json = match http_client
648 .info_spot_clearinghouse_state(&account_address)
649 .await
650 {
651 Ok(value) => value,
652 Err(e) => {
653 log::warn!("Settlement dispatch skipped: spot state fetch failed: {e}");
654 continue;
655 }
656 };
657 let spot_state: SpotClearinghouseState = match serde_json::from_value(spot_json) {
658 Ok(state) => state,
659 Err(e) => {
660 log::warn!("Settlement dispatch skipped: spot state parse failed: {e}");
661 continue;
662 }
663 };
664
665 let ts = clock.get_time_ns();
666 let fills = {
667 let mut guard = tracker.lock().expect(MUTEX_POISONED);
668 build_settlement_fills(&settlements, &spot_state, &mut guard, account_id, ts)
669 };
670
671 for fill in fills {
672 log::debug!(
673 "Dispatching outcome settlement fill: instrument={}, price={}, qty={}",
674 fill.instrument_id,
675 fill.last_px,
676 fill.last_qty,
677 );
678 emitter.send_fill_report(fill);
679 }
680 }
681 });
682
683 if let Some(previous) = self.settlement_poll_handle.replace(handle) {
684 previous.abort();
685 }
686
687 Ok(())
688 }
689
690 fn abort_pending_tasks(&self) {
691 self.pending_tasks.abort_all();
692 }
693}
694
695#[async_trait(?Send)]
696impl ExecutionClient for HyperliquidExecutionClient {
697 fn is_connected(&self) -> bool {
698 self.core.is_connected()
699 }
700
701 fn client_id(&self) -> ClientId {
702 self.core.client_id
703 }
704
705 fn account_id(&self) -> AccountId {
706 self.core.account_id
707 }
708
709 fn venue(&self) -> Venue {
710 *HYPERLIQUID_VENUE
711 }
712
713 fn socket_reconnect_registry(&self) -> Option<&SocketReconnectRegistry> {
714 Some(&self.socket_registry)
715 }
716
717 fn oms_type(&self) -> OmsType {
718 self.core.oms_type
719 }
720
721 fn get_account(&self) -> Option<AccountAny> {
722 self.core.cache().account_owned(&self.core.account_id)
723 }
724
725 fn generate_account_state(
726 &self,
727 balances: Vec<AccountBalance>,
728 margins: Vec<MarginBalance>,
729 reported: bool,
730 ts_event: UnixNanos,
731 info: Option<Params>,
732 ) -> anyhow::Result<()> {
733 self.emitter
734 .emit_account_state(balances, margins, reported, ts_event, info);
735 Ok(())
736 }
737
738 fn start(&mut self) -> anyhow::Result<()> {
739 if self.core.is_started() {
740 return Ok(());
741 }
742
743 let sender = get_exec_event_sender();
744 self.emitter.set_sender(sender);
745 self.core.set_started();
746
747 log::info!(
748 "Started: client_id={}, account_id={}, environment={:?}, vault_address={:?}, proxy_url={:?}",
749 self.core.client_id,
750 self.core.account_id,
751 self.config.environment,
752 self.config.vault_address,
753 self.config.proxy_url,
754 );
755
756 Ok(())
757 }
758
759 fn stop(&mut self) -> anyhow::Result<()> {
760 if self.core.is_stopped() {
761 return Ok(());
762 }
763
764 log::info!("Stopping Hyperliquid execution client");
765
766 if let Some(handle) = self.ws_stream_handle.take() {
767 handle.abort();
768 }
769
770 if let Some(handle) = self.settlement_poll_handle.take() {
771 handle.abort();
772 }
773
774 self.abort_pending_tasks();
775 self.ws_client.abort();
776
777 self.core.set_disconnected();
778 self.core.set_stopped();
779
780 log::info!("Hyperliquid execution client stopped");
781 Ok(())
782 }
783
784 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
785 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
786
787 if order.is_closed() {
788 log::warn!("Cannot submit closed order {}", order.client_order_id());
789 return Ok(());
790 }
791
792 if let Err(e) = self.validate_order_submission(&order) {
793 self.emitter
794 .emit_order_denied(&order, &format!("Validation failed: {e}"));
795 return Err(e);
796 }
797
798 let http_client = self.http_client.clone();
799 let symbol = order.instrument_id().symbol.inner();
800
801 let asset = match http_client.get_asset_index_for_symbol(symbol) {
803 Some(a) => a,
804 None => {
805 self.emitter
806 .emit_order_denied(&order, &format!("Asset index not found for {symbol}"));
807 return Ok(());
808 }
809 };
810
811 let price_decimals = http_client
813 .get_price_precision_for_symbol(symbol)
814 .unwrap_or(2);
815 let slippage_bps = self.resolve_slippage_bps(cmd.params.as_ref());
816 let mut hyperliquid_order = match order_to_hyperliquid_request_with_asset_and_cloid(
817 &order,
818 asset,
819 price_decimals,
820 self.config.normalize_prices,
821 slippage_bps,
822 None,
823 ) {
824 Ok(req) => req,
825 Err(e) => {
826 self.emitter
827 .emit_order_denied(&order, &format!("Order conversion failed: {e}"));
828 return Ok(());
829 }
830 };
831 let cloid = http_client.get_or_generate_client_order_id_cloid(order.client_order_id());
832 hyperliquid_order.cloid = Some(cloid);
833 if order.order_type() == OrderType::Market {
835 let instrument_id = order.instrument_id();
836 let cache = self.core.cache();
837 match cache.quote(&instrument_id) {
838 Some(quote) => {
839 let is_buy = order.order_side() == OrderSide::Buy;
840 hyperliquid_order.price =
841 derive_market_order_price(quote, is_buy, price_decimals, slippage_bps);
842 }
843 None => {
844 self.emitter.emit_order_denied(
845 &order,
846 &format!(
847 "No cached quote for {instrument_id}: \
848 subscribe to quote data before submitting market orders"
849 ),
850 );
851 return Ok(());
852 }
853 }
854 }
855
856 log::debug!(
857 "Submitting order: id={}, type={:?}, side={:?}, price={}, size={}, kind={:?}",
858 order.client_order_id(),
859 order.order_type(),
860 order.order_side(),
861 hyperliquid_order.price,
862 hyperliquid_order.size,
863 hyperliquid_order.kind,
864 );
865
866 let cloid = hyperliquid_order
869 .cloid
870 .expect("order conversion must set a CLOID");
871 self.http_client
872 .cache_client_order_id_cloid(order.client_order_id(), cloid);
873 self.ws_client
874 .cache_cloid_mapping(Ustr::from(&cloid.to_hex()), order.client_order_id());
875
876 self.register_order_context(&order);
877
878 self.emitter.emit_order_submitted(&order);
879
880 let emitter = self.emitter.clone();
881 let clock = self.clock;
882 let ws_client = self.ws_client.clone();
883 let cloid_hex = Ustr::from(&cloid.to_hex());
884 let dispatch_state = self.ws_dispatch_state.clone();
885
886 let builder = self.http_client.builder_attribution();
887
888 self.spawn_task("submit_order", async move {
889 let action = HyperliquidExecAction::Order {
890 orders: vec![hyperliquid_order],
891 grouping: HyperliquidExecGrouping::Na,
892 builder,
893 };
894 let rejection_route =
895 PostRejectionRoute::new(&emitter, &ws_client, &http_client, dispatch_state.clone());
896
897 match ws_client.post_action_exec(&http_client, &action).await {
898 Ok(response) => {
899 if response.is_ok() {
900 if let Some(inner_error) = extract_inner_error(&response) {
901 log::warn!("Order submission rejected by exchange: {inner_error}");
902 let ts = clock.get_time_ns();
903 rejection_route.emit_once(&order, &inner_error, ts, &cloid_hex);
904 } else {
905 log::debug!("Order submitted successfully: {response:?}");
906 }
907 } else {
908 let error_msg = extract_error_message(&response);
909 log::warn!("Order submission rejected by exchange: {error_msg}");
910 let ts = clock.get_time_ns();
911 rejection_route.emit_once(&order, &error_msg, ts, &cloid_hex);
912 }
913 }
914 Err(e) => {
915 log::error!("Order submission WebSocket post request failed: {e}");
919 }
920 }
921 rejection_route.resolve_without_post_rejection(&order, clock.get_time_ns(), &cloid_hex);
922
923 Ok(())
924 });
925
926 Ok(())
927 }
928
929 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
930 log::debug!(
931 "Submitting order list with {} orders",
932 cmd.order_list.client_order_ids.len()
933 );
934
935 let http_client = self.http_client.clone();
936 let slippage_bps = self.resolve_slippage_bps(cmd.params.as_ref());
937
938 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
939
940 let mut valid_orders = Vec::new();
941 let mut hyperliquid_orders = Vec::new();
942
943 for order in &orders {
944 match self.order_request(order, slippage_bps) {
945 Ok(request) => {
946 hyperliquid_orders.push(request);
947 valid_orders.push(order.clone());
948 }
949 Err(e) => {
950 self.emitter
951 .emit_order_denied(order, &format!("Order conversion failed: {e}"));
952 }
953 }
954 }
955
956 if valid_orders.is_empty() {
957 log::warn!("No valid orders to submit in order list");
958 return Ok(());
959 }
960
961 let grouping = determine_order_list_grouping(&valid_orders);
962 log::debug!("Order list grouping: {grouping:?}");
963 let (mut valid_orders, mut hyperliquid_orders) =
964 order_normal_tpsl_submission(valid_orders, hyperliquid_orders, grouping);
965
966 let submission_grouping = if grouping == HyperliquidExecGrouping::NormalTpsl {
967 let parent = valid_orders.remove(0);
968 let parent_request = hyperliquid_orders.remove(0);
969 let children = valid_orders
970 .drain(..)
971 .zip(hyperliquid_orders.drain(..))
972 .map(|(order, request)| StagedBracketChild { order, request })
973 .collect();
974 self.staged_brackets
975 .lock()
976 .expect(MUTEX_POISONED)
977 .stage(parent.client_order_id(), children);
978 valid_orders.push(parent);
979 hyperliquid_orders.push(parent_request);
980 HyperliquidExecGrouping::Na
981 } else {
982 grouping
983 };
984
985 for (order, request) in valid_orders.iter().zip(hyperliquid_orders.iter()) {
986 let cloid = request.cloid.expect("order conversion must set a CLOID");
987 self.http_client
988 .cache_client_order_id_cloid(order.client_order_id(), cloid);
989 self.ws_client
990 .cache_cloid_mapping(Ustr::from(&cloid.to_hex()), order.client_order_id());
991 self.register_order_context(order);
992 self.emitter.emit_order_submitted(order);
993 }
994
995 let emitter = self.emitter.clone();
996 let clock = self.clock;
997 let ws_client = self.ws_client.clone();
998 let dispatch_state = self.ws_dispatch_state.clone();
999 let staged_brackets = self.staged_brackets.clone();
1000 let builder = self.http_client.builder_attribution();
1001
1002 self.spawn_task("submit_order_list", async move {
1003 post_order_batch(
1004 "Order list",
1005 valid_orders,
1006 hyperliquid_orders,
1007 submission_grouping,
1008 builder,
1009 &emitter,
1010 &ws_client,
1011 &http_client,
1012 dispatch_state,
1013 staged_brackets,
1014 clock,
1015 )
1016 .await;
1017
1018 Ok(())
1019 });
1020
1021 Ok(())
1022 }
1023
1024 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1025 log::debug!("Modifying order: {cmd:?}");
1026
1027 let client_order_id = cmd.client_order_id;
1028 let venue_order_id = cmd
1029 .venue_order_id
1030 .or_else(|| self.core.cache().venue_order_id(&client_order_id).copied());
1031
1032 let order = match self.core.cache().order(&client_order_id).map(|o| o.clone()) {
1034 Some(o) => o,
1035 None => {
1036 let reason = "order not found in cache";
1037 log::warn!("Cannot modify order {client_order_id}: {reason}");
1038 self.emitter.emit_order_modify_rejected_event(
1039 cmd.strategy_id,
1040 cmd.instrument_id,
1041 client_order_id,
1042 venue_order_id,
1043 reason,
1044 self.clock.get_time_ns(),
1045 );
1046 return Ok(());
1047 }
1048 };
1049
1050 let http_client = self.http_client.clone();
1051 let symbol = cmd.instrument_id.symbol.inner();
1052 let should_normalize = self.config.normalize_prices;
1053 let slippage_bps = self.resolve_slippage_bps(cmd.params.as_ref());
1054 let modify_target = match http_client.unique_cached_client_order_id_cloid(&client_order_id)
1055 {
1056 Some(cloid) => HyperliquidExecModifyTarget::Cloid(cloid),
1057 None => {
1058 let Some(venue_order_id) = venue_order_id.as_ref() else {
1059 let reason = "venue_order_id or unique cached CLOID is required for modify";
1060 log::warn!("Cannot modify order {client_order_id}: {reason}");
1061 self.emitter.emit_order_modify_rejected_event(
1062 cmd.strategy_id,
1063 cmd.instrument_id,
1064 client_order_id,
1065 None,
1066 reason,
1067 self.clock.get_time_ns(),
1068 );
1069 return Ok(());
1070 };
1071
1072 match HyperliquidExecModifyTarget::from_venue_order_id(venue_order_id) {
1073 Ok(target) => target,
1074 Err(e) => {
1075 let reason =
1076 format!("Failed to parse venue_order_id '{venue_order_id}': {e}");
1077 log::warn!("{reason}");
1078 self.emitter.emit_order_modify_rejected_event(
1079 cmd.strategy_id,
1080 cmd.instrument_id,
1081 client_order_id,
1082 Some(*venue_order_id),
1083 &reason,
1084 self.clock.get_time_ns(),
1085 );
1086 return Ok(());
1087 }
1088 }
1089 }
1090 };
1091 let old_venue_order_id = venue_order_id.filter(|id| id.as_str().parse::<u64>().is_ok());
1092 if matches!(modify_target, HyperliquidExecModifyTarget::Cloid(_))
1093 && old_venue_order_id.is_none()
1094 {
1095 let reason = "cached venue_order_id is required for CLOID modify";
1096 log::warn!("Cannot modify order {client_order_id}: {reason}");
1097 self.emitter.emit_order_modify_rejected_event(
1098 cmd.strategy_id,
1099 cmd.instrument_id,
1100 client_order_id,
1101 venue_order_id,
1102 reason,
1103 self.clock.get_time_ns(),
1104 );
1105 return Ok(());
1106 }
1107
1108 let target_total_qty = cmd.quantity.unwrap_or(order.quantity());
1110 let filled_qty = order.filled_qty();
1111 if target_total_qty <= filled_qty {
1112 let reason =
1113 format!("modify quantity {target_total_qty} not greater than filled {filled_qty}",);
1114 log::warn!("Cannot modify order {}: {reason}", cmd.client_order_id);
1115
1116 self.emitter.emit_order_modify_rejected_event(
1117 cmd.strategy_id,
1118 cmd.instrument_id,
1119 client_order_id,
1120 venue_order_id,
1121 &reason,
1122 self.clock.get_time_ns(),
1123 );
1124 return Ok(());
1125 }
1126
1127 let quantity = target_total_qty - filled_qty;
1128 let price_decimals = http_client
1129 .get_price_precision_for_symbol(symbol)
1130 .unwrap_or(2);
1131 let asset = match http_client.get_asset_index_for_symbol(symbol) {
1132 Some(a) => a,
1133 None => {
1134 log::warn!(
1135 "Asset index not found for symbol {symbol}, ensure instruments are loaded",
1136 );
1137 return Ok(());
1138 }
1139 };
1140
1141 let mut hyperliquid_order = match order_to_hyperliquid_request_with_asset_and_cloid(
1144 &order,
1145 asset,
1146 price_decimals,
1147 should_normalize,
1148 slippage_bps,
1149 None,
1150 ) {
1151 Ok(mut req) => {
1152 if let Some(p) = cmd.price.or(order.price()) {
1154 let price_dec = p.as_decimal();
1155 req.price = if should_normalize {
1156 normalize_price(price_dec, price_decimals).normalize()
1157 } else {
1158 price_dec.normalize()
1159 };
1160 } else if let Some(tp) = cmd.trigger_price {
1161 let is_buy = order.order_side() == OrderSide::Buy;
1164 let base = tp.as_decimal().normalize();
1165 let derived = derive_limit_from_trigger(base, is_buy, slippage_bps);
1166 let sig_rounded = round_to_sig_figs(derived, 5);
1167 req.price =
1168 clamp_price_to_precision(sig_rounded, price_decimals, is_buy).normalize();
1169 }
1170 req.size = quantity.as_decimal().normalize();
1173
1174 if let (Some(tp), HyperliquidExecOrderKind::Trigger { trigger }) =
1176 (cmd.trigger_price, &mut req.kind)
1177 {
1178 let tp_dec = tp.as_decimal();
1179 trigger.trigger_px = if should_normalize {
1180 normalize_price(tp_dec, price_decimals).normalize()
1181 } else {
1182 tp_dec.normalize()
1183 };
1184 }
1185
1186 req
1187 }
1188 Err(e) => {
1189 log::warn!("Order conversion failed for modify: {e}");
1190 return Ok(());
1191 }
1192 };
1193 let cached_cloid_before_modify = http_client.cached_client_order_id_cloid(&client_order_id);
1194 let cloid = http_client.get_or_generate_client_order_id_cloid(order.client_order_id());
1195 let generated_modify_cloid = cached_cloid_before_modify
1196 .is_none()
1197 .then_some((client_order_id, cloid));
1198 hyperliquid_order.cloid = Some(cloid);
1199
1200 let dispatch_state = self.ws_dispatch_state.clone();
1201 let ws_client = self.ws_client.clone();
1202
1203 if let Some(cloid) = hyperliquid_order.cloid {
1204 http_client.cache_client_order_id_cloid(client_order_id, cloid);
1205 ws_client.cache_cloid_mapping(Ustr::from(&cloid.to_hex()), client_order_id);
1206 }
1207
1208 let modify_generation = old_venue_order_id.map(|old_venue_order_id| {
1211 let generation = dispatch_state.mark_pending_modify(
1212 client_order_id,
1213 old_venue_order_id,
1214 target_total_qty,
1215 );
1216 dispatch_state.stash_modify_request(client_order_id, hyperliquid_order.clone());
1218 generation
1219 });
1220
1221 self.spawn_task("modify_order", async move {
1222 let action = HyperliquidExecAction::Modify {
1223 modify: HyperliquidExecModifyOrderRequest {
1224 oid: modify_target,
1225 order: hyperliquid_order,
1226 },
1227 };
1228
1229 match ws_client.post_action_exec(&http_client, &action).await {
1230 Ok(response) => {
1231 if response.is_ok() {
1232 if let Some(inner_error) = extract_inner_error(&response) {
1233 log::warn!("Order modification rejected by exchange: {inner_error}");
1234
1235 if let Some(generation) = modify_generation {
1236 dispatch_state
1237 .clear_modify_generation(&client_order_id, generation);
1238 }
1239 remove_generated_modify_cloid(
1240 &http_client,
1241 &ws_client,
1242 generated_modify_cloid,
1243 );
1244 } else {
1245 log::debug!("Order modified successfully: {response:?}");
1246 }
1247 } else {
1248 let error_msg = extract_error_message(&response);
1249 log::warn!("Order modification rejected by exchange: {error_msg}");
1250
1251 if let Some(generation) = modify_generation {
1252 dispatch_state.clear_modify_generation(&client_order_id, generation);
1253 }
1254 remove_generated_modify_cloid(
1255 &http_client,
1256 &ws_client,
1257 generated_modify_cloid,
1258 );
1259 }
1260 }
1261 Err(e) => {
1262 if e.is_transport_error() {
1263 log::warn!(
1265 "Order modification transport failure for {client_order_id}: {e}; \
1266 awaiting WS reconciliation",
1267 );
1268 } else {
1269 log::warn!("Order modification WebSocket post request failed: {e}");
1270
1271 if let Some(generation) = modify_generation {
1272 dispatch_state.clear_modify_generation(&client_order_id, generation);
1273 }
1274 remove_generated_modify_cloid(
1275 &http_client,
1276 &ws_client,
1277 generated_modify_cloid,
1278 );
1279 }
1280 }
1281 }
1282
1283 Ok(())
1284 });
1285
1286 Ok(())
1287 }
1288
1289 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1290 log::debug!("Cancelling order: {cmd:?}");
1291
1292 if let Some(order) = self
1293 .staged_brackets
1294 .lock()
1295 .expect(MUTEX_POISONED)
1296 .cancel_child(&cmd.client_order_id)
1297 {
1298 self.emitter
1299 .emit_order_canceled(&order, None, self.clock.get_time_ns());
1300 return Ok(());
1301 }
1302
1303 let http_client = self.http_client.clone();
1304 let emitter = self.emitter.clone();
1305 let clock = self.clock;
1306 let client_order_id = cmd.client_order_id;
1307 let strategy_id = cmd.strategy_id;
1308 let instrument_id = cmd.instrument_id;
1309 let venue_order_id = cmd.venue_order_id;
1310 let symbol = cmd.instrument_id.symbol.inner();
1311 let ws_client = self.ws_client.clone();
1312 let fast = can_fast_cancel_order(
1313 self.core
1314 .cache()
1315 .order(&client_order_id)
1316 .as_ref()
1317 .map(|order| order.order_type()),
1318 )
1319 .then_some(true);
1320
1321 self.spawn_task("cancel_order", async move {
1322 let asset = match http_client.get_asset_index_for_symbol(symbol) {
1323 Some(a) => a,
1324 None => {
1325 log::warn!(
1326 "Local cancel validation failed for {client_order_id}: Asset index not found for symbol {symbol}"
1327 );
1328 return Ok(());
1329 }
1330 };
1331
1332 let action =
1333 if let Some(cloid) = http_client.cached_client_order_id_cloid(&client_order_id) {
1334 HyperliquidExecAction::CancelByCloid {
1335 cancels: vec![HyperliquidExecCancelByCloidRequest { asset, cloid }],
1336 fast,
1337 }
1338 } else if let Some(venue_order_id) = venue_order_id {
1339 match venue_order_id.as_str().parse::<u64>() {
1340 Ok(oid) => HyperliquidExecAction::Cancel {
1341 cancels: vec![HyperliquidExecCancelOrderRequest { asset, oid }],
1342 fast,
1343 },
1344 Err(_) => {
1345 log::warn!(
1346 "Local cancel validation failed for {client_order_id}: Invalid venue order ID format"
1347 );
1348 return Ok(());
1349 }
1350 }
1351 } else {
1352 let cloid = http_client.get_or_generate_client_order_id_cloid(client_order_id);
1353 HyperliquidExecAction::CancelByCloid {
1354 cancels: vec![HyperliquidExecCancelByCloidRequest { asset, cloid }],
1355 fast,
1356 }
1357 };
1358
1359 match ws_client.post_action_exec(&http_client, &action).await {
1360 Ok(response) => {
1361 if response.is_ok() {
1362 if let Some(inner_error) = extract_inner_error(&response) {
1363 emitter.emit_order_cancel_rejected_event(
1364 strategy_id,
1365 instrument_id,
1366 client_order_id,
1367 venue_order_id,
1368 &inner_error,
1369 clock.get_time_ns(),
1370 );
1371 } else {
1372 log::debug!("Order cancelled successfully: {response:?}");
1373 }
1374 } else {
1375 let error_msg = extract_error_message(&response);
1376 log::warn!(
1377 "Cancel failed without per-order result for {client_order_id}, awaiting WS reconciliation: {error_msg}"
1378 );
1379 }
1380 }
1381 Err(e) => {
1382 if e.is_transport_error() {
1383 log::warn!(
1384 "Cancel transport failure for {client_order_id}: {e}; \
1385 awaiting WS reconciliation",
1386 );
1387 } else {
1388 log::warn!(
1389 "Ambiguous cancel failure for {client_order_id}, awaiting WS reconciliation: {e}"
1390 );
1391 }
1392 }
1393 }
1394
1395 Ok(())
1396 });
1397
1398 Ok(())
1399 }
1400
1401 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
1402 log::debug!("Cancelling all orders: {cmd:?}");
1403
1404 let cache = self.core.cache();
1405 let open_orders = cache.orders_open(
1406 Some(&self.core.venue),
1407 Some(&cmd.instrument_id),
1408 None,
1409 None,
1410 Some(cmd.order_side),
1411 );
1412
1413 if open_orders.is_empty() {
1414 log::debug!("No open orders to cancel for {:?}", cmd.instrument_id);
1415 return Ok(());
1416 }
1417
1418 let symbol = cmd.instrument_id.symbol.inner();
1419 let instrument_id = cmd.instrument_id;
1420 let strategy_id = cmd.strategy_id;
1421 let entries: Vec<CancelEntry> = open_orders
1422 .iter()
1423 .map(|o| CancelEntry {
1424 strategy_id,
1425 instrument_id,
1426 client_order_id: o.client_order_id(),
1427 venue_order_id: o.venue_order_id(),
1428 symbol,
1429 fast: can_fast_cancel_order(Some(o.order_type())),
1430 })
1431 .collect();
1432
1433 let http_client = self.http_client.clone();
1434 let emitter = self.emitter.clone();
1435 let clock = self.clock;
1436 let ws_client = self.ws_client.clone();
1437
1438 self.spawn_task("cancel_all_orders", async move {
1439 let asset = match http_client.get_asset_index_for_symbol(symbol) {
1440 Some(a) => a,
1441 None => {
1442 log::warn!(
1443 "Local cancel-all validation failed: Asset index not found for symbol {symbol}"
1444 );
1445 return Ok(());
1446 }
1447 };
1448
1449 let mut cancel_dispatch = CancelDispatch::new();
1450
1451 for entry in &entries {
1452 cancel_dispatch.push(entry, asset, &http_client);
1453 }
1454
1455 if cancel_dispatch.is_empty() {
1456 return Ok(());
1457 }
1458
1459 submit_cancel_dispatch(
1460 "Cancel-all",
1461 cancel_dispatch,
1462 &ws_client,
1463 &http_client,
1464 &emitter,
1465 clock,
1466 )
1467 .await;
1468
1469 Ok(())
1470 });
1471
1472 Ok(())
1473 }
1474
1475 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1476 log::debug!("Batch cancelling orders: {cmd:?}");
1477
1478 if cmd.cancels.is_empty() {
1479 log::debug!("No orders to cancel in batch");
1480 return Ok(());
1481 }
1482
1483 let cache = self.core.cache();
1484 let entries: Vec<CancelEntry> = cmd
1485 .cancels
1486 .iter()
1487 .map(|c| CancelEntry {
1488 strategy_id: c.strategy_id,
1489 instrument_id: c.instrument_id,
1490 client_order_id: c.client_order_id,
1491 venue_order_id: c.venue_order_id,
1492 symbol: c.instrument_id.symbol.inner(),
1493 fast: can_fast_cancel_order(
1494 cache
1495 .order(&c.client_order_id)
1496 .as_ref()
1497 .map(|order| order.order_type()),
1498 ),
1499 })
1500 .collect();
1501
1502 let http_client = self.http_client.clone();
1503 let emitter = self.emitter.clone();
1504 let clock = self.clock;
1505 let ws_client = self.ws_client.clone();
1506
1507 self.spawn_task("batch_cancel_orders", async move {
1508 let mut cancel_dispatch = CancelDispatch::new();
1509
1510 for entry in &entries {
1511 let asset = match http_client.get_asset_index_for_symbol(entry.symbol) {
1512 Some(a) => a,
1513 None => {
1514 log::warn!(
1515 "Local batch cancel validation failed for {}: Asset index not found for symbol {}",
1516 entry.client_order_id,
1517 entry.symbol,
1518 );
1519 continue;
1520 }
1521 };
1522
1523 cancel_dispatch.push(entry, asset, &http_client);
1524 }
1525
1526 if cancel_dispatch.is_empty() {
1527 log::warn!("No valid cancel requests in batch");
1528 return Ok(());
1529 }
1530
1531 submit_cancel_dispatch(
1532 "Batch cancel",
1533 cancel_dispatch,
1534 &ws_client,
1535 &http_client,
1536 &emitter,
1537 clock,
1538 )
1539 .await;
1540
1541 Ok(())
1542 });
1543
1544 Ok(())
1545 }
1546
1547 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
1548 let http_client = self.http_client.clone();
1549 let account_address = self.get_account_address()?;
1550 let emitter = self.emitter.clone();
1551 let clock = self.clock;
1552
1553 self.spawn_task("query_account", async move {
1554 let perp_json = http_client
1555 .info_clearinghouse_state(&account_address)
1556 .await
1557 .context("failed to fetch clearinghouse state")?;
1558
1559 let perp_state: ClearinghouseState = serde_json::from_value(perp_json)
1560 .context("failed to deserialize clearinghouse state")?;
1561
1562 let spot_json = http_client
1563 .info_spot_clearinghouse_state(&account_address)
1564 .await
1565 .context("failed to fetch spot clearinghouse state")?;
1566 let spot_state: SpotClearinghouseState = serde_json::from_value(spot_json)
1567 .context("failed to deserialize spot clearinghouse state")?;
1568
1569 let (balances, margins) =
1570 parse_combined_account_balances_and_margins(&perp_state, &spot_state)
1571 .context("failed to parse combined account balances and margins")?;
1572 let ts_event = clock.get_time_ns();
1573 emitter.emit_account_state(balances, margins, true, ts_event, None);
1574
1575 Ok(())
1576 });
1577
1578 Ok(())
1579 }
1580
1581 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
1582 log::debug!("Querying order: {cmd:?}");
1583
1584 let client_order_id = cmd.client_order_id;
1585 let venue_order_id = match cmd.venue_order_id {
1586 Some(voi) => Some(voi),
1587 None => self.core.cache().venue_order_id(&client_order_id).copied(),
1588 };
1589
1590 let account_address = self.get_account_address()?;
1591 let http_client = self.http_client.clone();
1592 let emitter = self.emitter.clone();
1593 let dispatch_state = self.ws_dispatch_state.clone();
1594 let clock = self.clock;
1595
1596 self.spawn_task("query_order", async move {
1597 match http_client
1602 .request_order_status_report_by_client_order_id(&account_address, &client_order_id)
1603 .await
1604 {
1605 Ok(Some(report)) => {
1606 promote_replacement_from_query(
1607 &report,
1608 &dispatch_state,
1609 &emitter,
1610 clock.get_time_ns(),
1611 );
1612 log::debug!("Queried order status for {client_order_id}");
1613 emitter.send_order_status_report(report);
1614 return Ok(());
1615 }
1616 Ok(None) => {}
1617 Err(e) => {
1618 log::warn!(
1619 "Failed to query order status for {client_order_id}: {e}; falling back to oid lookup"
1620 );
1621 }
1622 }
1623
1624 let Some(venue_order_id) = venue_order_id else {
1625 log::debug!("No order status report found for {client_order_id}");
1626 return Ok(());
1627 };
1628
1629 let oid: u64 = match venue_order_id.as_str().parse() {
1630 Ok(oid) => oid,
1631 Err(e) => {
1632 log::warn!("Failed to parse venue order ID {venue_order_id}: {e}");
1633 return Ok(());
1634 }
1635 };
1636
1637 match http_client
1638 .request_order_status_report(&account_address, oid)
1639 .await
1640 {
1641 Ok(Some(mut report)) => {
1642 if is_inflight_modify_old_leg_cancel(
1643 &dispatch_state,
1644 &client_order_id,
1645 &report,
1646 ) {
1647 log::debug!(
1648 "Suppressing stale old-leg Canceled for {client_order_id}: modify in flight"
1649 );
1650 } else {
1651 attach_known_client_order_id(&mut report, client_order_id);
1652 log::debug!("Queried order status for oid {oid}");
1653 emitter.send_order_status_report(report);
1654 }
1655 }
1656 Ok(None) => {
1657 log::debug!("No order status report found for oid {oid}");
1658 }
1659 Err(e) => {
1660 log::warn!("Failed to query order status for oid {oid}: {e}");
1661 }
1662 }
1663
1664 Ok(())
1665 });
1666
1667 Ok(())
1668 }
1669
1670 async fn connect(&mut self) -> anyhow::Result<()> {
1671 if self.core.is_connected() {
1672 return Ok(());
1673 }
1674
1675 log::info!("Connecting Hyperliquid execution client");
1676
1677 self.ensure_instruments_initialized_async().await?;
1679 let ready_bracket_parents = self.restore_staged_brackets();
1680
1681 self.start_ws_stream().await?;
1683
1684 let post_ws = async {
1686 self.refresh_account_state().await?;
1687 self.await_account_registered(30.0).await?;
1688
1689 Ok::<(), anyhow::Error>(())
1690 };
1691
1692 if let Err(e) = post_ws.await {
1693 log::warn!("Connect failed after WS started, tearing down: {e}");
1694 let _ = self.ws_client.disconnect().await;
1695 self.abort_pending_tasks();
1696 return Err(e);
1697 }
1698
1699 for parent_id in ready_bracket_parents {
1700 if let Some(children) = self
1701 .staged_brackets
1702 .lock()
1703 .expect(MUTEX_POISONED)
1704 .activate(&parent_id)
1705 {
1706 spawn_staged_children(
1707 children,
1708 &self.emitter,
1709 &self.ws_client,
1710 &self.http_client,
1711 self.ws_dispatch_state.clone(),
1712 self.staged_brackets.clone(),
1713 self.http_client.builder_attribution(),
1714 self.clock,
1715 );
1716 }
1717 }
1718
1719 if let Err(e) = self.start_outcome_settlement_poll() {
1720 log::warn!("Outcome settlement polling not started: {e}");
1721 }
1722
1723 self.core.set_connected();
1724
1725 log::info!("Connected: client_id={}", self.core.client_id);
1726 Ok(())
1727 }
1728
1729 async fn disconnect(&mut self) -> anyhow::Result<()> {
1730 if self.core.is_disconnected() {
1731 return Ok(());
1732 }
1733
1734 log::info!("Disconnecting Hyperliquid execution client");
1735
1736 self.ws_client.disconnect().await?;
1738
1739 if let Some(handle) = self.settlement_poll_handle.take() {
1740 handle.abort();
1741 }
1742
1743 self.abort_pending_tasks();
1745
1746 self.core.set_disconnected();
1747
1748 log::info!("Disconnected: client_id={}", self.core.client_id);
1749 Ok(())
1750 }
1751
1752 async fn generate_order_status_report(
1753 &self,
1754 cmd: &GenerateOrderStatusReport,
1755 ) -> anyhow::Result<Option<OrderStatusReport>> {
1756 let account_address = self.get_account_address()?;
1757
1758 if cmd.venue_order_id.is_none() && cmd.client_order_id.is_none() {
1759 log::warn!(
1760 "Cannot generate order status report without venue_order_id or client_order_id"
1761 );
1762 return Ok(None);
1763 }
1764
1765 if let Some(client_order_id) = &cmd.client_order_id {
1769 match self
1770 .http_client
1771 .request_order_status_report_by_client_order_id(&account_address, client_order_id)
1772 .await
1773 {
1774 Ok(Some(report)) => {
1775 promote_replacement_from_query(
1776 &report,
1777 &self.ws_dispatch_state,
1778 &self.emitter,
1779 self.clock.get_time_ns(),
1780 );
1781 log::debug!("Generated order status report for {client_order_id}");
1782 return Ok(Some(report));
1783 }
1784 Ok(None) => {}
1785 Err(e) => {
1786 log::warn!(
1787 "Failed to generate order status report for {client_order_id}: {e}; \
1788 falling back to oid lookup"
1789 );
1790 }
1791 }
1792 }
1793
1794 let oid = match &cmd.venue_order_id {
1795 Some(venue_order_id) => venue_order_id
1796 .as_str()
1797 .parse::<u64>()
1798 .context("failed to parse venue_order_id as oid")?,
1799 None => match &cmd.client_order_id {
1800 Some(client_order_id) => {
1801 let cached_oid: Option<u64> = self
1802 .core
1803 .cache()
1804 .venue_order_id(client_order_id)
1805 .and_then(|v| v.as_str().parse::<u64>().ok());
1806
1807 match cached_oid {
1808 Some(oid) => oid,
1809 None => {
1810 log::debug!("No order status report found for {client_order_id}");
1811 return Ok(None);
1812 }
1813 }
1814 }
1815 None => unreachable!("cmd must carry at least one identifier"),
1816 },
1817 };
1818
1819 let mut report = self
1820 .http_client
1821 .request_order_status_report(&account_address, oid)
1822 .await
1823 .context("failed to generate order status report")?;
1824
1825 if let Some(report) = &report
1826 && let Some(client_order_id) = &cmd.client_order_id
1827 && is_inflight_modify_old_leg_cancel(&self.ws_dispatch_state, client_order_id, report)
1828 {
1829 log::debug!(
1830 "Suppressing stale old-leg Canceled for {client_order_id}: modify in flight"
1831 );
1832 return Ok(None);
1833 }
1834
1835 if let Some(report) = &mut report
1836 && let Some(client_order_id) = cmd.client_order_id
1837 {
1838 attach_known_client_order_id(report, client_order_id);
1839 }
1840
1841 if report.is_some() {
1842 log::debug!("Generated order status report for oid {oid}");
1843 } else {
1844 log::debug!("No order status report found for oid {oid}");
1845 }
1846 Ok(report)
1847 }
1848
1849 async fn generate_order_status_reports(
1850 &self,
1851 cmd: &GenerateOrderStatusReports,
1852 ) -> anyhow::Result<Vec<OrderStatusReport>> {
1853 let account_address = self.get_account_address()?;
1854
1855 let reports = self
1856 .http_client
1857 .request_order_status_reports(&account_address, cmd.instrument_id)
1858 .await
1859 .context("failed to generate order status reports")?;
1860
1861 let reports = filter_order_status_reports_for_command(reports, cmd);
1862
1863 log::debug!("Generated {} order status reports", reports.len());
1864 Ok(reports)
1865 }
1866
1867 async fn generate_fill_reports(
1868 &self,
1869 cmd: GenerateFillReports,
1870 ) -> anyhow::Result<Vec<FillReport>> {
1871 let account_address = self.get_account_address()?;
1872
1873 let reports = self
1874 .http_client
1875 .request_fill_reports(&account_address, cmd.instrument_id)
1876 .await
1877 .context("failed to generate fill reports")?;
1878
1879 let reports = if let (Some(start), Some(end)) = (cmd.start, cmd.end) {
1881 reports
1882 .into_iter()
1883 .filter(|r| r.ts_event >= start && r.ts_event <= end)
1884 .collect()
1885 } else if let Some(start) = cmd.start {
1886 reports
1887 .into_iter()
1888 .filter(|r| r.ts_event >= start)
1889 .collect()
1890 } else if let Some(end) = cmd.end {
1891 reports.into_iter().filter(|r| r.ts_event <= end).collect()
1892 } else {
1893 reports
1894 };
1895
1896 log::debug!("Generated {} fill reports", reports.len());
1897 Ok(reports)
1898 }
1899
1900 async fn generate_position_status_reports(
1901 &self,
1902 cmd: &GeneratePositionStatusReports,
1903 ) -> anyhow::Result<Vec<PositionStatusReport>> {
1904 let account_address = self.get_account_address()?;
1905
1906 let reports = self
1908 .http_client
1909 .request_position_status_reports(&account_address, cmd.instrument_id)
1910 .await
1911 .context("failed to generate position status reports")?;
1912
1913 log::debug!("Generated {} position status reports", reports.len());
1914 Ok(reports)
1915 }
1916
1917 async fn generate_mass_status(
1918 &self,
1919 lookback_mins: Option<u64>,
1920 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
1921 let ts_init = self.clock.get_time_ns();
1922
1923 let order_cmd = GenerateOrderStatusReports::new(
1924 UUID4::new(),
1925 ts_init,
1926 true, None,
1928 None,
1929 None,
1930 None,
1931 None,
1932 );
1933 let fill_cmd =
1934 GenerateFillReports::new(UUID4::new(), ts_init, None, None, None, None, None, None);
1935 let position_cmd =
1936 GeneratePositionStatusReports::new(UUID4::new(), ts_init, None, None, None, None, None);
1937
1938 let mut order_reports = self.generate_order_status_reports(&order_cmd).await?;
1939 let mut fill_reports = self.generate_fill_reports(fill_cmd).await?;
1940 let position_reports = self.generate_position_status_reports(&position_cmd).await?;
1941
1942 if let Some(mins) = lookback_mins {
1945 let cutoff_ns = ts_init
1946 .as_u64()
1947 .saturating_sub(mins.saturating_mul(60).saturating_mul(1_000_000_000));
1948 let cutoff = UnixNanos::from(cutoff_ns);
1949
1950 fill_reports.retain(|r| r.ts_event >= cutoff);
1951 }
1952
1953 if !fill_reports.is_empty() {
1954 let account_address = self.get_account_address()?;
1955 let filled_order_ids: ahash::AHashSet<_> = fill_reports
1956 .iter()
1957 .map(|report| report.venue_order_id)
1958 .collect();
1959 let open_order_ids: ahash::AHashSet<_> = order_reports
1960 .iter()
1961 .map(|report| report.venue_order_id)
1962 .collect();
1963 let mut historical_reports = self
1964 .http_client
1965 .request_historical_order_status_reports(&account_address, None)
1966 .await
1967 .context("failed to generate historical order status reports")?;
1968 historical_reports.retain(|report| {
1969 filled_order_ids.contains(&report.venue_order_id)
1970 && !open_order_ids.contains(&report.venue_order_id)
1971 });
1972 order_reports.extend(historical_reports);
1973 }
1974
1975 let mut mass_status = ExecutionMassStatus::new(
1976 self.core.client_id,
1977 self.core.account_id,
1978 self.core.venue,
1979 ts_init,
1980 None,
1981 );
1982 mass_status.add_order_reports(order_reports);
1983 mass_status.add_fill_reports(fill_reports);
1984 mass_status.add_position_reports(position_reports);
1985
1986 log::info!(
1987 "Generated mass status: {} orders, {} fills, {} positions",
1988 mass_status.order_reports().len(),
1989 mass_status.fill_reports().len(),
1990 mass_status.position_reports().len(),
1991 );
1992
1993 Ok(Some(mass_status))
1994 }
1995}
1996
1997impl HyperliquidExecutionClient {
1998 async fn start_ws_stream(&mut self) -> anyhow::Result<()> {
1999 if self.ws_stream_handle.is_some() {
2000 return Ok(());
2001 }
2002
2003 let subscription_address = self.get_account_address()?;
2005
2006 let mut ws_client = self.ws_client.clone();
2007
2008 let instruments = self
2009 .http_client
2010 .request_instruments()
2011 .await
2012 .unwrap_or_default();
2013
2014 for instrument in instruments {
2015 ws_client.cache_instrument(instrument);
2016 }
2017
2018 ws_client.connect().await?;
2020 ws_client
2021 .subscribe_order_updates(&subscription_address)
2022 .await?;
2023 ws_client
2024 .subscribe_user_events(&subscription_address)
2025 .await?;
2026 log::debug!("Subscribed to Hyperliquid execution updates for {subscription_address}");
2027
2028 if let Some(handle) = ws_client.take_task_handle() {
2030 self.ws_client.set_task_handle(handle);
2031 }
2032
2033 let emitter = self.emitter.clone();
2034 let dispatch_state = self.ws_dispatch_state.clone();
2035 let staged_brackets = self.staged_brackets.clone();
2036 let http_client = self.http_client.clone();
2037 let builder = self.http_client.builder_attribution();
2038 let clock = self.clock;
2039 let runtime = get_runtime();
2040 let handle = runtime.spawn(async move {
2041 let mut pending_filled_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
2052
2053 loop {
2054 let event = ws_client.next_event().await;
2055
2056 match event {
2057 Some(msg) => match msg {
2058 NautilusWsMessage::ExecutionReports(reports) => {
2059 for report in reports {
2060 let staged_parent_fill = match &report {
2061 ExecutionReport::Fill(report) => report.client_order_id,
2062 ExecutionReport::Order(_) => None,
2063 };
2064
2065 let staged_parent_terminal = match &report {
2066 ExecutionReport::Order(report)
2067 if matches!(
2068 report.order_status,
2069 OrderStatus::Canceled
2070 | OrderStatus::Rejected
2071 | OrderStatus::Expired
2072 ) =>
2073 {
2074 report.client_order_id.map(|client_order_id| {
2075 (client_order_id, report.ts_last)
2076 })
2077 }
2078 _ => None,
2079 };
2080
2081 let active_child_terminal = match &report {
2082 ExecutionReport::Order(report)
2083 if matches!(
2084 report.order_status,
2085 OrderStatus::Filled
2086 | OrderStatus::Canceled
2087 | OrderStatus::Rejected
2088 | OrderStatus::Expired
2089 ) =>
2090 {
2091 report.client_order_id
2092 }
2093 ExecutionReport::Fill(report) => {
2094 report.client_order_id.filter(|client_order_id| {
2095 let Some(context) =
2096 dispatch_state.lookup_context(client_order_id)
2097 else {
2098 return false;
2099 };
2100 let previous = dispatch_state
2101 .previous_filled_qty(client_order_id)
2102 .unwrap_or_else(|| {
2103 Quantity::zero(report.last_qty.precision)
2104 });
2105 previous + report.last_qty >= context.quantity
2106 })
2107 }
2108 _ => None,
2109 };
2110
2111 let active_child_fill = match &report {
2112 ExecutionReport::Fill(report) => {
2113 report.client_order_id.and_then(|client_order_id| {
2114 dispatch_state.lookup_context(&client_order_id).map(
2115 |context| {
2116 (
2117 client_order_id,
2118 dispatch_state
2119 .previous_filled_qty(&client_order_id)
2120 .unwrap_or_else(|| {
2121 Quantity::zero(
2122 report.last_qty.precision,
2123 )
2124 }),
2125 context.quantity,
2126 )
2127 },
2128 )
2129 })
2130 }
2131 ExecutionReport::Order(_) => None,
2132 };
2133
2134 if let Some((cid, oid, order)) = handle_execution_report(
2135 report,
2136 &dispatch_state,
2137 &emitter,
2138 &ws_client,
2139 &http_client,
2140 &mut pending_filled_cloids,
2141 clock.get_time_ns(),
2142 ) {
2143 spawn_corrective_reduce(
2144 &ws_client,
2145 &http_client,
2146 &dispatch_state,
2147 cid,
2148 oid,
2149 order,
2150 );
2151 }
2152
2153 if let Some(parent_id) = staged_parent_fill
2154 && let Some(children) = staged_brackets
2155 .lock()
2156 .expect(MUTEX_POISONED)
2157 .activate(&parent_id)
2158 {
2159 spawn_staged_children(
2160 children,
2161 &emitter,
2162 &ws_client,
2163 &http_client,
2164 dispatch_state.clone(),
2165 staged_brackets.clone(),
2166 builder.clone(),
2167 clock,
2168 );
2169 }
2170
2171 if let Some((parent_id, ts_event)) = staged_parent_terminal {
2172 let children = staged_brackets
2173 .lock()
2174 .expect(MUTEX_POISONED)
2175 .cancel_for_parent(&parent_id);
2176
2177 for child in children {
2178 emitter.emit_order_canceled(&child, None, ts_event);
2179 }
2180 }
2181
2182 if let Some((client_order_id, previous, quantity)) =
2183 active_child_fill
2184 && let Some(cumulative) =
2185 dispatch_state.previous_filled_qty(&client_order_id)
2186 && cumulative > previous
2187 && cumulative < quantity
2188 {
2189 let sibling = staged_brackets
2190 .lock()
2191 .expect(MUTEX_POISONED)
2192 .active_sibling(&client_order_id);
2193
2194 if let Some(sibling) = sibling {
2195 spawn_active_sibling_resize(
2196 sibling,
2197 quantity - cumulative,
2198 &emitter,
2199 &ws_client,
2200 &http_client,
2201 &dispatch_state,
2202 );
2203 }
2204 }
2205
2206 if let Some(client_order_id) = active_child_terminal {
2207 let sibling = staged_brackets
2208 .lock()
2209 .expect(MUTEX_POISONED)
2210 .take_active_sibling(&client_order_id);
2211
2212 if let Some(sibling) = sibling {
2213 spawn_active_sibling_cancel(
2214 sibling,
2215 &emitter,
2216 &ws_client,
2217 &http_client,
2218 &dispatch_state,
2219 );
2220 }
2221 }
2222 }
2223 }
2224 NautilusWsMessage::Reconnected => {
2225 log::info!("WebSocket reconnected");
2226 }
2227 NautilusWsMessage::Error(e) => {
2228 log::warn!("WebSocket error: {e}");
2229 }
2230 NautilusWsMessage::Trades(_)
2232 | NautilusWsMessage::Quote(_)
2233 | NautilusWsMessage::Deltas(_)
2234 | NautilusWsMessage::Depth10(_)
2235 | NautilusWsMessage::Candle(_)
2236 | NautilusWsMessage::MarkPrice(_)
2237 | NautilusWsMessage::IndexPrice(_)
2238 | NautilusWsMessage::FundingRate(_)
2239 | NautilusWsMessage::CustomData(_) => {}
2240 },
2241 None => {
2242 log::debug!("WebSocket next_event returned None, stream closed");
2243 break;
2244 }
2245 }
2246 }
2247 });
2248
2249 self.ws_stream_handle = Some(handle);
2250 log::debug!("Hyperliquid WebSocket execution stream started");
2251 Ok(())
2252 }
2253}
2254
2255fn filter_order_status_reports_for_command(
2256 reports: Vec<OrderStatusReport>,
2257 cmd: &GenerateOrderStatusReports,
2258) -> Vec<OrderStatusReport> {
2259 let reports = if cmd.open_only {
2260 reports
2261 .into_iter()
2262 .filter(|r| r.order_status.is_open())
2263 .collect()
2264 } else {
2265 reports
2266 };
2267
2268 match (cmd.start, cmd.end) {
2269 (Some(start), Some(end)) => reports
2270 .into_iter()
2271 .filter(|r| r.ts_last >= start && r.ts_last <= end)
2272 .collect(),
2273 (Some(start), None) => reports.into_iter().filter(|r| r.ts_last >= start).collect(),
2274 (None, Some(end)) => reports.into_iter().filter(|r| r.ts_last <= end).collect(),
2275 (None, None) => reports,
2276 }
2277}
2278
2279fn attach_known_client_order_id(report: &mut OrderStatusReport, client_order_id: ClientOrderId) {
2280 if report.client_order_id.is_none() {
2281 report.client_order_id = Some(client_order_id);
2282 }
2283}
2284
2285fn is_inflight_modify_old_leg_cancel(
2289 dispatch_state: &WsDispatchState,
2290 client_order_id: &ClientOrderId,
2291 report: &OrderStatusReport,
2292) -> bool {
2293 report.order_status == OrderStatus::Canceled
2294 && dispatch_state.pending_modify_contains_old(client_order_id, report.venue_order_id)
2295}
2296
2297fn remove_generated_modify_cloid(
2298 http_client: &HyperliquidHttpClient,
2299 ws_client: &HyperliquidWebSocketClient,
2300 generated_modify_cloid: Option<(ClientOrderId, Cloid)>,
2301) {
2302 let Some((client_order_id, cloid)) = generated_modify_cloid else {
2303 return;
2304 };
2305
2306 if http_client.cached_client_order_id_cloid(&client_order_id) != Some(cloid) {
2307 return;
2308 }
2309
2310 let cloid_hex = Ustr::from(&cloid.to_hex());
2311 ws_client.remove_cloid_mapping(&cloid_hex);
2312 http_client.remove_client_order_id_cloid(&client_order_id);
2313}
2314
2315#[derive(Clone)]
2316struct CancelEntry {
2317 strategy_id: StrategyId,
2318 instrument_id: InstrumentId,
2319 client_order_id: ClientOrderId,
2320 venue_order_id: Option<VenueOrderId>,
2321 symbol: Ustr,
2322 fast: bool,
2323}
2324
2325struct CancelDispatch {
2326 cloid_requests: Vec<(HyperliquidExecCancelByCloidRequest, CancelEntry)>,
2327 oid_requests: Vec<(HyperliquidExecCancelOrderRequest, CancelEntry)>,
2328}
2329
2330impl CancelDispatch {
2331 fn new() -> Self {
2332 Self {
2333 cloid_requests: Vec::new(),
2334 oid_requests: Vec::new(),
2335 }
2336 }
2337
2338 fn is_empty(&self) -> bool {
2339 self.cloid_requests.is_empty() && self.oid_requests.is_empty()
2340 }
2341
2342 fn push(&mut self, entry: &CancelEntry, asset: u32, http_client: &HyperliquidHttpClient) {
2343 if let Some(cloid) = http_client.cached_client_order_id_cloid(&entry.client_order_id) {
2344 self.cloid_requests.push((
2345 HyperliquidExecCancelByCloidRequest { asset, cloid },
2346 entry.clone(),
2347 ));
2348 } else if let Some(venue_order_id) = entry.venue_order_id {
2349 match venue_order_id.as_str().parse::<u64>() {
2350 Ok(oid) => {
2351 self.oid_requests.push((
2352 HyperliquidExecCancelOrderRequest { asset, oid },
2353 entry.clone(),
2354 ));
2355 }
2356 Err(_) => {
2357 log::warn!(
2358 "Local cancel validation failed for {}: Invalid venue order ID format",
2359 entry.client_order_id,
2360 );
2361 }
2362 }
2363 } else {
2364 let cloid = http_client.get_or_generate_client_order_id_cloid(entry.client_order_id);
2365 self.cloid_requests.push((
2366 HyperliquidExecCancelByCloidRequest { asset, cloid },
2367 entry.clone(),
2368 ));
2369 }
2370 }
2371}
2372
2373async fn submit_cancel_dispatch(
2374 label: &str,
2375 dispatch: CancelDispatch,
2376 ws_client: &HyperliquidWebSocketClient,
2377 http_client: &HyperliquidHttpClient,
2378 emitter: &ExecutionEventEmitter,
2379 clock: &'static AtomicTime,
2380) {
2381 let CancelDispatch {
2382 cloid_requests,
2383 oid_requests,
2384 } = dispatch;
2385
2386 let (fast_cloid_requests, fast_cloid_entries, cloid_requests, cloid_entries) =
2387 split_fast_cancel_requests(cloid_requests);
2388
2389 if !fast_cloid_requests.is_empty() {
2390 let action = HyperliquidExecAction::CancelByCloid {
2391 cancels: fast_cloid_requests,
2392 fast: Some(true),
2393 };
2394 submit_cancel_action(
2395 label,
2396 action,
2397 &fast_cloid_entries,
2398 ws_client,
2399 http_client,
2400 emitter,
2401 clock,
2402 )
2403 .await;
2404 }
2405
2406 if !cloid_requests.is_empty() {
2407 let action = HyperliquidExecAction::CancelByCloid {
2408 cancels: cloid_requests,
2409 fast: None,
2410 };
2411 submit_cancel_action(
2412 label,
2413 action,
2414 &cloid_entries,
2415 ws_client,
2416 http_client,
2417 emitter,
2418 clock,
2419 )
2420 .await;
2421 }
2422
2423 let (fast_oid_requests, fast_oid_entries, oid_requests, oid_entries) =
2424 split_fast_cancel_requests(oid_requests);
2425
2426 if !fast_oid_requests.is_empty() {
2427 let action = HyperliquidExecAction::Cancel {
2428 cancels: fast_oid_requests,
2429 fast: Some(true),
2430 };
2431 submit_cancel_action(
2432 label,
2433 action,
2434 &fast_oid_entries,
2435 ws_client,
2436 http_client,
2437 emitter,
2438 clock,
2439 )
2440 .await;
2441 }
2442
2443 if !oid_requests.is_empty() {
2444 let action = HyperliquidExecAction::Cancel {
2445 cancels: oid_requests,
2446 fast: None,
2447 };
2448 submit_cancel_action(
2449 label,
2450 action,
2451 &oid_entries,
2452 ws_client,
2453 http_client,
2454 emitter,
2455 clock,
2456 )
2457 .await;
2458 }
2459}
2460
2461fn split_fast_cancel_requests<T>(
2462 requests: Vec<(T, CancelEntry)>,
2463) -> (Vec<T>, Vec<CancelEntry>, Vec<T>, Vec<CancelEntry>) {
2464 let mut fast_requests = Vec::new();
2465 let mut fast_entries = Vec::new();
2466 let mut requests_without_fast = Vec::new();
2467 let mut entries_without_fast = Vec::new();
2468
2469 for (request, entry) in requests {
2470 if entry.fast {
2471 fast_requests.push(request);
2472 fast_entries.push(entry);
2473 } else {
2474 requests_without_fast.push(request);
2475 entries_without_fast.push(entry);
2476 }
2477 }
2478
2479 (
2480 fast_requests,
2481 fast_entries,
2482 requests_without_fast,
2483 entries_without_fast,
2484 )
2485}
2486
2487async fn submit_cancel_action(
2488 label: &str,
2489 action: HyperliquidExecAction,
2490 sent_entries: &[CancelEntry],
2491 ws_client: &HyperliquidWebSocketClient,
2492 http_client: &HyperliquidHttpClient,
2493 emitter: &ExecutionEventEmitter,
2494 clock: &'static AtomicTime,
2495) {
2496 match ws_client.post_action_exec(http_client, &action).await {
2497 Ok(response) => {
2498 if response.is_ok() {
2499 let inner_errors = extract_inner_errors(&response);
2500 let ts = clock.get_time_ns();
2501
2502 if inner_errors.is_empty() {
2503 log::debug!("{label} submitted successfully: {response:?}");
2504 } else if let Some(reason) = cancel_status_count_mismatch_reason(
2505 label,
2506 sent_entries.len(),
2507 inner_errors.len(),
2508 ) {
2509 log::warn!("{reason}");
2510 } else {
2511 for (i, entry) in sent_entries.iter().enumerate() {
2512 if let Some(Some(error_msg)) = inner_errors.get(i) {
2513 log::warn!(
2514 "Cancel for {} rejected by exchange: {error_msg}",
2515 entry.client_order_id,
2516 );
2517 emitter.emit_order_cancel_rejected_event(
2518 entry.strategy_id,
2519 entry.instrument_id,
2520 entry.client_order_id,
2521 entry.venue_order_id,
2522 error_msg,
2523 ts,
2524 );
2525 }
2526 }
2527 }
2528 } else {
2529 let error_msg = extract_error_message(&response);
2530 log::warn!(
2531 "{label} failed without per-order results, awaiting WS reconciliation: {error_msg}"
2532 );
2533 }
2534 }
2535 Err(e) => {
2536 if e.is_transport_error() {
2537 log::warn!("{label} transport failure: {e}; awaiting WS reconciliation");
2538 } else {
2539 log::warn!("{label} ambiguous failure, awaiting WS reconciliation: {e}");
2540 }
2541 }
2542 }
2543}
2544
2545fn register_order_context_into(state: &WsDispatchState, order: &OrderAny) {
2554 let context = OrderContext::from(order);
2555 if context.is_quote_quantity {
2556 return;
2557 }
2558
2559 state.register_context(context);
2560 state.mark_submission_pending(context.identity.client_order_id);
2561}
2562
2563fn order_normal_tpsl_submission(
2564 orders: Vec<OrderAny>,
2565 requests: Vec<HyperliquidExecPlaceOrderRequest>,
2566 grouping: HyperliquidExecGrouping,
2567) -> (Vec<OrderAny>, Vec<HyperliquidExecPlaceOrderRequest>) {
2568 if grouping != HyperliquidExecGrouping::NormalTpsl {
2569 return (orders, requests);
2570 }
2571
2572 let mut pairs: Vec<_> = orders.into_iter().zip(requests).collect();
2573 pairs.sort_by_key(|(order, request)| {
2574 if !order.is_reduce_only() {
2575 0
2576 } else if matches!(
2577 &request.kind,
2578 HyperliquidExecOrderKind::Trigger { trigger }
2579 if trigger.tpsl == HyperliquidExecTpSl::Sl
2580 ) {
2581 2
2582 } else {
2583 1
2584 }
2585 });
2586
2587 pairs.into_iter().unzip()
2588}
2589
2590pub fn validate_order_for_hyperliquid(order: &OrderAny) -> anyhow::Result<()> {
2599 let instrument_id = order.instrument_id();
2600 let symbol = instrument_id.symbol.as_str();
2601 let product_type = HyperliquidProductType::from_symbol(symbol).map_err(|_| {
2602 anyhow::anyhow!(
2603 "Unsupported instrument symbol format for Hyperliquid: {symbol} \
2604 (expected -PERP, -SPOT, or HIP-4 outcome `{{N}}-{{YES|NO}}-OUTCOME`)"
2605 )
2606 })?;
2607
2608 match order.order_type() {
2609 OrderType::Market
2610 | OrderType::Limit
2611 | OrderType::StopMarket
2612 | OrderType::StopLimit
2613 | OrderType::MarketIfTouched
2614 | OrderType::LimitIfTouched => {}
2615 _ => anyhow::bail!(
2616 "Unsupported order type for Hyperliquid: {:?}",
2617 order.order_type()
2618 ),
2619 }
2620
2621 if product_type == HyperliquidProductType::Outcome {
2624 if order.is_reduce_only() {
2625 anyhow::bail!("Reduce-only is not supported for Hyperliquid HIP-4 outcomes: {symbol}");
2626 }
2627
2628 if !matches!(order.order_type(), OrderType::Market | OrderType::Limit) {
2629 anyhow::bail!(
2630 "Trigger order types are not supported for Hyperliquid HIP-4 outcomes: \
2631 {symbol} (received {:?})",
2632 order.order_type()
2633 );
2634 }
2635 }
2636
2637 if matches!(
2638 order.order_type(),
2639 OrderType::StopMarket
2640 | OrderType::StopLimit
2641 | OrderType::MarketIfTouched
2642 | OrderType::LimitIfTouched
2643 ) && order.trigger_price().is_none()
2644 {
2645 anyhow::bail!(
2646 "Conditional orders require a trigger price for Hyperliquid: {:?}",
2647 order.order_type()
2648 );
2649 }
2650
2651 if matches!(
2652 order.order_type(),
2653 OrderType::Limit | OrderType::StopLimit | OrderType::LimitIfTouched
2654 ) && order.price().is_none()
2655 {
2656 anyhow::bail!(
2657 "Limit orders require a limit price for Hyperliquid: {:?}",
2658 order.order_type()
2659 );
2660 }
2661
2662 Ok(())
2663}
2664
2665fn can_fast_cancel_order(order_type: Option<OrderType>) -> bool {
2666 matches!(order_type, Some(OrderType::Market | OrderType::Limit))
2667}
2668
2669fn cancel_status_count_mismatch_reason(
2670 label: &str,
2671 expected_count: usize,
2672 actual_count: usize,
2673) -> Option<String> {
2674 (actual_count != 0 && actual_count != expected_count).then(|| {
2675 format!(
2676 "{label} response status count mismatch: expected {expected_count}, received {actual_count}"
2677 )
2678 })
2679}
2680
2681#[expect(clippy::too_many_arguments)]
2682async fn post_order_batch(
2683 label: &str,
2684 orders: Vec<OrderAny>,
2685 requests: Vec<HyperliquidExecPlaceOrderRequest>,
2686 grouping: HyperliquidExecGrouping,
2687 builder: Option<crate::http::models::HyperliquidExecBuilderFee>,
2688 emitter: &ExecutionEventEmitter,
2689 ws_client: &HyperliquidWebSocketClient,
2690 http_client: &HyperliquidHttpClient,
2691 dispatch_state: Arc<WsDispatchState>,
2692 staged_brackets: Arc<Mutex<StagedBracketState>>,
2693 clock: &'static AtomicTime,
2694) {
2695 let cloid_hexes: Vec<Ustr> = requests
2696 .iter()
2697 .map(|request| {
2698 Ustr::from(
2699 &request
2700 .cloid
2701 .expect("order conversion must set a CLOID")
2702 .to_hex(),
2703 )
2704 })
2705 .collect();
2706 let action = HyperliquidExecAction::Order {
2707 orders: requests,
2708 grouping,
2709 builder,
2710 };
2711 let rejection_route = PostRejectionRoute::with_staged_brackets(
2712 emitter,
2713 ws_client,
2714 http_client,
2715 dispatch_state,
2716 staged_brackets,
2717 );
2718
2719 match ws_client.post_action_exec(http_client, &action).await {
2720 Ok(response) if response.is_ok() => {
2721 let inner_errors = extract_inner_errors(&response);
2722 let ts = clock.get_time_ns();
2723
2724 if inner_errors.len() == orders.len() {
2725 for ((order, cloid_hex), error) in orders
2726 .iter()
2727 .zip(cloid_hexes.iter())
2728 .zip(inner_errors.iter())
2729 {
2730 if let Some(error_msg) = error {
2731 log::warn!(
2732 "Order {} rejected by exchange: {error_msg}",
2733 order.client_order_id(),
2734 );
2735 rejection_route.emit_once(order, error_msg, ts, cloid_hex);
2736 }
2737 }
2738 } else if orders.len() > 1
2739 && inner_errors.len() == 1
2740 && let Some(error_msg) = inner_errors[0].as_ref()
2741 {
2742 log::warn!("{label} rejected by deterministic whole-batch validation: {error_msg}",);
2743 for (order, cloid_hex) in orders.iter().zip(cloid_hexes.iter()) {
2744 rejection_route.emit_once(order, error_msg, ts, cloid_hex);
2745 }
2746 } else if !inner_errors.is_empty() {
2747 log::warn!(
2748 "{label} returned {} statuses for {} orders; preserving unresolved identities \
2749 for WebSocket or startup reconciliation",
2750 inner_errors.len(),
2751 orders.len(),
2752 );
2753 } else {
2754 log::debug!("{label} submitted successfully: {response:?}");
2755 }
2756 }
2757 Ok(response) => {
2758 let error_msg = extract_error_message(&response);
2759 log::warn!("{label} submission rejected by exchange: {error_msg}");
2760 let ts = clock.get_time_ns();
2761
2762 for (order, cloid_hex) in orders.iter().zip(cloid_hexes.iter()) {
2763 rejection_route.emit_once(order, &error_msg, ts, cloid_hex);
2764 }
2765 }
2766 Err(e) => {
2767 log::error!("{label} WebSocket post request failed: {e}");
2770 }
2771 }
2772
2773 let ts = clock.get_time_ns();
2774 for (order, cloid_hex) in orders.iter().zip(cloid_hexes.iter()) {
2775 rejection_route.resolve_without_post_rejection(order, ts, cloid_hex);
2776 }
2777}
2778
2779#[expect(clippy::too_many_arguments)]
2780fn spawn_staged_children(
2781 children: Vec<StagedBracketChild>,
2782 emitter: &ExecutionEventEmitter,
2783 ws_client: &HyperliquidWebSocketClient,
2784 http_client: &HyperliquidHttpClient,
2785 dispatch_state: Arc<WsDispatchState>,
2786 staged_brackets: Arc<Mutex<StagedBracketState>>,
2787 builder: Option<crate::http::models::HyperliquidExecBuilderFee>,
2788 clock: &'static AtomicTime,
2789) {
2790 let (orders, requests): (Vec<_>, Vec<_>) = children
2791 .into_iter()
2792 .map(|child| (child.order, child.request))
2793 .unzip();
2794
2795 for (order, request) in orders.iter().zip(requests.iter()) {
2796 let cloid = request.cloid.expect("order conversion must set a CLOID");
2797 http_client.cache_client_order_id_cloid(order.client_order_id(), cloid);
2798 ws_client.cache_cloid_mapping(Ustr::from(&cloid.to_hex()), order.client_order_id());
2799 register_order_context_into(&dispatch_state, order);
2800 emitter.emit_order_submitted(order);
2801 }
2802
2803 let emitter = emitter.clone();
2804 let ws_client = ws_client.clone();
2805 let http_client = http_client.clone();
2806
2807 get_runtime().spawn(async move {
2808 post_order_batch(
2809 "Bracket child batch",
2810 orders,
2811 requests,
2812 HyperliquidExecGrouping::Na,
2813 builder,
2814 &emitter,
2815 &ws_client,
2816 &http_client,
2817 dispatch_state,
2818 staged_brackets,
2819 clock,
2820 )
2821 .await;
2822 });
2823}
2824
2825fn spawn_active_sibling_cancel(
2826 sibling: StagedBracketChild,
2827 emitter: &ExecutionEventEmitter,
2828 ws_client: &HyperliquidWebSocketClient,
2829 http_client: &HyperliquidHttpClient,
2830 dispatch_state: &WsDispatchState,
2831) {
2832 let client_order_id = sibling.order.client_order_id();
2833 let Some(cloid) = sibling.request.cloid else {
2834 log::error!("Cannot cancel OUO sibling {client_order_id}: missing CLOID");
2835 return;
2836 };
2837 let venue_order_id = dispatch_state.cached_venue_order_id(&client_order_id);
2838 let action = HyperliquidExecAction::CancelByCloid {
2839 cancels: vec![HyperliquidExecCancelByCloidRequest {
2840 asset: sibling.request.asset,
2841 cloid,
2842 }],
2843 fast: can_fast_cancel_order(Some(sibling.order.order_type())).then_some(true),
2844 };
2845 let emitter = emitter.clone();
2846 let ws_client = ws_client.clone();
2847 let http_client = http_client.clone();
2848
2849 get_runtime().spawn(async move {
2850 match ws_client.post_action_exec(&http_client, &action).await {
2851 Ok(response) if response.is_ok() => {
2852 if let Some(error) = extract_inner_error(&response) {
2853 emitter.emit_order_cancel_rejected(
2854 &sibling.order,
2855 venue_order_id,
2856 &error,
2857 get_atomic_clock_realtime().get_time_ns(),
2858 );
2859 }
2860 }
2861 Ok(response) => {
2862 log::warn!(
2863 "OUO sibling cancel for {client_order_id} returned an ambiguous response; \
2864 awaiting WebSocket or startup reconciliation: {}",
2865 extract_error_message(&response),
2866 );
2867 }
2868 Err(e) => {
2869 log::warn!(
2870 "OUO sibling cancel for {client_order_id} failed; awaiting WebSocket or \
2871 startup reconciliation: {e}",
2872 );
2873 }
2874 }
2875 });
2876}
2877
2878fn spawn_active_sibling_resize(
2879 sibling: StagedBracketChild,
2880 target_total_qty: Quantity,
2881 emitter: &ExecutionEventEmitter,
2882 ws_client: &HyperliquidWebSocketClient,
2883 http_client: &HyperliquidHttpClient,
2884 dispatch_state: &Arc<WsDispatchState>,
2885) {
2886 let client_order_id = sibling.order.client_order_id();
2887 let Some(old_venue_order_id) = dispatch_state.cached_venue_order_id(&client_order_id) else {
2888 log::warn!(
2889 "Cannot resize OUO sibling {client_order_id}: venue order ID not known; awaiting \
2890 WebSocket or startup reconciliation",
2891 );
2892 return;
2893 };
2894 let filled_qty = dispatch_state
2895 .previous_filled_qty(&client_order_id)
2896 .unwrap_or_else(|| Quantity::zero(target_total_qty.precision));
2897 let Some(order) = build_ouo_resize_request(&sibling, target_total_qty, filled_qty) else {
2898 spawn_active_sibling_cancel(sibling, emitter, ws_client, http_client, dispatch_state);
2899 return;
2900 };
2901 let Some(cloid) = order.cloid else {
2902 log::error!("Cannot resize OUO sibling {client_order_id}: missing CLOID");
2903 return;
2904 };
2905
2906 let generation =
2907 dispatch_state.mark_pending_modify(client_order_id, old_venue_order_id, target_total_qty);
2908 dispatch_state.stash_modify_request(client_order_id, order.clone());
2909 let action = HyperliquidExecAction::Modify {
2910 modify: HyperliquidExecModifyOrderRequest {
2911 oid: HyperliquidExecModifyTarget::Cloid(cloid),
2912 order,
2913 },
2914 };
2915 let ws_client = ws_client.clone();
2916 let http_client = http_client.clone();
2917 let dispatch_state = dispatch_state.clone();
2918
2919 get_runtime().spawn(async move {
2920 match ws_client.post_action_exec(&http_client, &action).await {
2921 Ok(response) if response.is_ok() && extract_inner_error(&response).is_none() => {
2922 log::debug!("OUO sibling resize submitted for {client_order_id}");
2923 }
2924 Ok(response) => {
2925 dispatch_state.clear_modify_generation(&client_order_id, generation);
2926 log::warn!(
2927 "OUO sibling resize for {client_order_id} rejected: {}",
2928 extract_inner_error(&response)
2929 .unwrap_or_else(|| extract_error_message(&response)),
2930 );
2931 }
2932 Err(e) if e.is_transport_error() => {
2933 log::warn!(
2934 "OUO sibling resize transport failure for {client_order_id}: {e}; awaiting \
2935 WebSocket or startup reconciliation",
2936 );
2937 }
2938 Err(e) => {
2939 dispatch_state.clear_modify_generation(&client_order_id, generation);
2940 log::warn!("OUO sibling resize failed for {client_order_id}: {e}");
2941 }
2942 }
2943 });
2944}
2945
2946fn build_ouo_resize_request(
2947 sibling: &StagedBracketChild,
2948 target_total_qty: Quantity,
2949 filled_qty: Quantity,
2950) -> Option<HyperliquidExecPlaceOrderRequest> {
2951 if target_total_qty <= filled_qty {
2952 return None;
2953 }
2954
2955 let mut request = sibling.request.clone();
2956 request.size = (target_total_qty - filled_qty).as_decimal().normalize();
2957 Some(request)
2958}
2959
2960struct PostRejectionRoute {
2961 emitter: ExecutionEventEmitter,
2962 ws_client: HyperliquidWebSocketClient,
2963 http_client: HyperliquidHttpClient,
2964 dispatch_state: Arc<WsDispatchState>,
2965 staged_brackets: Arc<Mutex<StagedBracketState>>,
2966}
2967
2968impl PostRejectionRoute {
2969 fn new(
2970 emitter: &ExecutionEventEmitter,
2971 ws_client: &HyperliquidWebSocketClient,
2972 http_client: &HyperliquidHttpClient,
2973 dispatch_state: Arc<WsDispatchState>,
2974 ) -> Self {
2975 Self {
2976 emitter: emitter.clone(),
2977 ws_client: ws_client.clone(),
2978 http_client: http_client.clone(),
2979 dispatch_state,
2980 staged_brackets: Arc::new(Mutex::new(StagedBracketState::default())),
2981 }
2982 }
2983
2984 fn with_staged_brackets(
2985 emitter: &ExecutionEventEmitter,
2986 ws_client: &HyperliquidWebSocketClient,
2987 http_client: &HyperliquidHttpClient,
2988 dispatch_state: Arc<WsDispatchState>,
2989 staged_brackets: Arc<Mutex<StagedBracketState>>,
2990 ) -> Self {
2991 Self {
2992 emitter: emitter.clone(),
2993 ws_client: ws_client.clone(),
2994 http_client: http_client.clone(),
2995 dispatch_state,
2996 staged_brackets,
2997 }
2998 }
2999
3000 fn emit_once(
3001 &self,
3002 order: &OrderAny,
3003 reason: &str,
3004 ts_event: UnixNanos,
3005 cloid_hex: &Ustr,
3006 ) -> bool {
3007 let client_order_id = order.client_order_id();
3008 let _ = self.dispatch_state.resolve_submission(&client_order_id);
3009
3010 if !self.dispatch_state.insert_filled(client_order_id) {
3011 log::debug!(
3012 "Skipping duplicate post rejection for terminal order {client_order_id}: {reason}",
3013 );
3014 self.ws_client.remove_cloid_mapping(cloid_hex);
3015 self.http_client
3016 .remove_client_order_id_cloid(&client_order_id);
3017 return false;
3018 }
3019
3020 if reason.contains(HYPERLIQUID_BUILDER_FEE_NOT_APPROVED) {
3021 log::warn!(
3022 "Builder fee not approved: complete the one-time 0% builder approval \
3023 (signed by the master wallet). See: {HYPERLIQUID_BUILDER_APPROVAL_DOCS_URL}",
3024 );
3025 }
3026
3027 let normalized_reason = reason.to_lowercase();
3028 let due_post_only = order.is_post_only()
3029 && (normalized_reason.contains(&HYPERLIQUID_POST_ONLY_WOULD_MATCH.to_lowercase())
3030 || normalized_reason.contains("post-only order would have immediately matched"));
3031 self.emitter
3032 .emit_order_rejected(order, reason, ts_event, due_post_only);
3033 let active_sibling = self
3034 .staged_brackets
3035 .lock()
3036 .expect(MUTEX_POISONED)
3037 .take_active_sibling(&client_order_id);
3038
3039 if let Some(sibling) = active_sibling {
3040 spawn_active_sibling_cancel(
3041 sibling,
3042 &self.emitter,
3043 &self.ws_client,
3044 &self.http_client,
3045 &self.dispatch_state,
3046 );
3047 }
3048 let staged_children = self
3049 .staged_brackets
3050 .lock()
3051 .expect(MUTEX_POISONED)
3052 .cancel_for_parent(&client_order_id);
3053
3054 for child in staged_children {
3055 self.emitter.emit_order_canceled(&child, None, ts_event);
3056 }
3057 self.dispatch_state.insert_terminal_cloid(*cloid_hex);
3058 self.dispatch_state.cleanup_terminal(&client_order_id);
3059 self.ws_client.remove_cloid_mapping(cloid_hex);
3060 self.http_client
3061 .remove_client_order_id_cloid(&client_order_id);
3062
3063 true
3064 }
3065
3066 fn resolve_without_post_rejection(
3067 &self,
3068 order: &OrderAny,
3069 ts_init: UnixNanos,
3070 cloid_hex: &Ustr,
3071 ) {
3072 let client_order_id = order.client_order_id();
3073 let Some(report) = self.dispatch_state.resolve_submission(&client_order_id) else {
3074 return;
3075 };
3076 let is_terminal = !report.order_status.is_open();
3077 let outcome = dispatch_order_event(&report, &self.dispatch_state, &self.emitter, ts_init);
3078
3079 if outcome == DispatchOutcome::External {
3080 self.emitter.send_order_status_report(report);
3081 }
3082
3083 if is_terminal && outcome != DispatchOutcome::Skip {
3084 self.ws_client.remove_cloid_mapping(cloid_hex);
3085 self.http_client
3086 .remove_client_order_id_cloid(&client_order_id);
3087 }
3088 }
3089}
3090
3091fn handle_execution_report(
3098 report: ExecutionReport,
3099 dispatch_state: &WsDispatchState,
3100 emitter: &ExecutionEventEmitter,
3101 ws_client: &HyperliquidWebSocketClient,
3102 http_client: &HyperliquidHttpClient,
3103 pending_filled_cloids: &mut FifoCache<ClientOrderId, 10_000>,
3104 ts_init: UnixNanos,
3105) -> Option<(ClientOrderId, u64, HyperliquidExecPlaceOrderRequest)> {
3106 match report {
3107 ExecutionReport::Order(order_report) => {
3108 let is_filled_marker = matches!(order_report.order_status, OrderStatus::Filled);
3109 let is_open = order_report.order_status.is_open();
3110 let client_order_id = order_report.client_order_id;
3111
3112 let outcome = dispatch_order_event(&order_report, dispatch_state, emitter, ts_init);
3113
3114 if outcome == DispatchOutcome::External {
3115 emitter.send_order_status_report(order_report);
3116 }
3117
3118 if let Some(id) = client_order_id
3130 && !is_open
3131 {
3132 match outcome {
3133 DispatchOutcome::Skip => {}
3134 DispatchOutcome::Tracked if is_filled_marker => {
3135 pending_filled_cloids.add(id);
3136 }
3137 DispatchOutcome::Tracked | DispatchOutcome::External => {
3138 remove_cloid_mapping_for_client_order_id(ws_client, http_client, &id);
3139 }
3140 }
3141 }
3142
3143 client_order_id.and_then(|id| {
3146 dispatch_state
3147 .take_corrective(&id)
3148 .map(|(oid, order)| (id, oid, order))
3149 })
3150 }
3151 ExecutionReport::Fill(fill_report) => {
3152 let client_order_id = fill_report.client_order_id;
3153
3154 let outcome = dispatch_order_fill(&fill_report, dispatch_state, emitter, ts_init);
3155
3156 if outcome == DispatchOutcome::External {
3157 emitter.send_fill_report(fill_report);
3158 }
3159
3160 if let Some(id) = client_order_id
3163 && pending_filled_cloids.contains(&id)
3164 && dispatch_state.buffered_fill_count(&id) == 0
3165 {
3166 pending_filled_cloids.remove(&id);
3167 remove_cloid_mapping_for_client_order_id(ws_client, http_client, &id);
3168 }
3169
3170 client_order_id.and_then(|id| {
3172 dispatch_state
3173 .take_corrective(&id)
3174 .map(|(oid, order)| (id, oid, order))
3175 })
3176 }
3177 }
3178}
3179
3180fn spawn_corrective_reduce(
3188 ws_client: &HyperliquidWebSocketClient,
3189 http_client: &HyperliquidHttpClient,
3190 dispatch_state: &Arc<WsDispatchState>,
3191 client_order_id: ClientOrderId,
3192 oid: u64,
3193 order: HyperliquidExecPlaceOrderRequest,
3194) {
3195 let ws_client = ws_client.clone();
3196 let http_client = http_client.clone();
3197 let dispatch_state = dispatch_state.clone();
3198
3199 get_runtime().spawn(async move {
3200 let action = HyperliquidExecAction::Modify {
3201 modify: HyperliquidExecModifyOrderRequest {
3202 oid: oid.into(),
3203 order,
3204 },
3205 };
3206
3207 let keep_marker = match ws_client.post_action_exec(&http_client, &action).await {
3208 Ok(resp) if resp.is_ok() && extract_inner_error(&resp).is_none() => {
3209 log::debug!("Corrective reduce acknowledged for {client_order_id} on oid {oid}");
3210 true
3211 }
3212 Ok(resp) => {
3213 let reason =
3214 extract_inner_error(&resp).unwrap_or_else(|| extract_error_message(&resp));
3215 log::warn!(
3216 "Corrective reduce rejected for {client_order_id} on oid {oid}: {reason}"
3217 );
3218 false
3219 }
3220 Err(e) if e.is_transport_error() => {
3221 log::warn!(
3222 "Corrective reduce transport failure for {client_order_id} on oid {oid}: \
3223 {e}; awaiting WS reconciliation",
3224 );
3225 true
3226 }
3227 Err(e) => {
3228 log::warn!("Corrective reduce failed for {client_order_id} on oid {oid}: {e}");
3229 false
3230 }
3231 };
3232
3233 if !keep_marker {
3234 dispatch_state.clear_pending_modify(&client_order_id);
3235 }
3236 });
3237}
3238
3239fn remove_cloid_mapping_for_client_order_id(
3240 ws_client: &HyperliquidWebSocketClient,
3241 http_client: &HyperliquidHttpClient,
3242 client_order_id: &ClientOrderId,
3243) {
3244 let generated_cloid = Cloid::from_client_order_id(*client_order_id);
3245
3246 if let Some(cloid) = http_client.remove_client_order_id_cloid(client_order_id) {
3247 ws_client.remove_cloid_mapping(&Ustr::from(&cloid.to_hex()));
3248 if cloid == generated_cloid {
3249 return;
3250 }
3251 }
3252
3253 ws_client.remove_cloid_mapping(&Ustr::from(&generated_cloid.to_hex()));
3254}
3255
3256use crate::common::parse::determine_order_list_grouping;
3257
3258#[cfg(test)]
3259mod tests {
3260 use std::sync::Arc;
3261
3262 use nautilus_common::messages::{ExecutionEvent, execution::GenerateOrderStatusReports};
3263 use nautilus_core::{UUID4, UnixNanos, time::get_atomic_clock_realtime};
3264 use nautilus_live::{
3265 ExecutionEventEmitter,
3266 execution::context::{OrderContext, OrderIdentity},
3267 };
3268 use nautilus_model::{
3269 enums::{
3270 AccountType, ContingencyType, LiquiditySide, OrderSide, OrderStatus, OrderType,
3271 TimeInForce, TriggerType,
3272 },
3273 events::OrderEventAny,
3274 identifiers::{
3275 AccountId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, VenueOrderId,
3276 },
3277 orders::{Order, OrderAny, limit::LimitOrder, stop_market::StopMarketOrder},
3278 reports::{FillReport, OrderStatusReport},
3279 types::{Currency, Money, Price, Quantity},
3280 };
3281 use nautilus_network::websocket::TransportBackend;
3282 use rstest::rstest;
3283 use rust_decimal::Decimal;
3284 use ustr::Ustr;
3285
3286 use super::{
3287 CancelEntry, ExecutionReport, FifoCache, HyperliquidHttpClient, HyperliquidWebSocketClient,
3288 PostRejectionRoute, StagedBracketChild, StagedBracketState, WsDispatchState,
3289 attach_known_client_order_id, build_ouo_resize_request, can_fast_cancel_order,
3290 determine_order_list_grouping, filter_order_status_reports_for_command,
3291 handle_execution_report, register_order_context_into, split_fast_cancel_requests,
3292 validate_order_for_hyperliquid,
3293 };
3294 use crate::{
3295 common::enums::HyperliquidEnvironment,
3296 http::models::{
3297 Cloid, HyperliquidExecGrouping, HyperliquidExecLimitParams, HyperliquidExecOrderKind,
3298 HyperliquidExecPlaceOrderRequest, HyperliquidExecTif,
3299 },
3300 };
3301
3302 const TEST_INSTRUMENT_ID: &str = "BTC-USD-PERP.HYPERLIQUID";
3303
3304 fn test_emitter() -> (
3305 ExecutionEventEmitter,
3306 tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3307 ) {
3308 let clock = get_atomic_clock_realtime();
3309 let mut emitter = ExecutionEventEmitter::new(
3310 clock,
3311 TraderId::from("TESTER-001"),
3312 AccountId::from("HYPERLIQUID-001"),
3313 AccountType::Margin,
3314 None,
3315 );
3316 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
3317 emitter.set_sender(tx);
3318 (emitter, rx)
3319 }
3320
3321 fn drain_events(
3322 rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3323 ) -> Vec<ExecutionEvent> {
3324 let mut out = Vec::new();
3325 while let Ok(e) = rx.try_recv() {
3326 out.push(e);
3327 }
3328 out
3329 }
3330
3331 fn make_ws_client() -> HyperliquidWebSocketClient {
3332 HyperliquidWebSocketClient::new(
3336 Some("wss://test.invalid".to_string()),
3337 HyperliquidEnvironment::Testnet,
3338 None,
3339 TransportBackend::default(),
3340 None,
3341 )
3342 }
3343
3344 fn make_http_client() -> HyperliquidHttpClient {
3345 HyperliquidHttpClient::new(HyperliquidEnvironment::Testnet, 1, None).unwrap()
3346 }
3347
3348 fn test_context(client_order_id: ClientOrderId) -> OrderContext {
3351 OrderContext {
3352 identity: OrderIdentity {
3353 client_order_id,
3354 strategy_id: StrategyId::from("S-001"),
3355 instrument_id: InstrumentId::from(TEST_INSTRUMENT_ID),
3356 order_side: OrderSide::Buy,
3357 order_type: OrderType::Limit,
3358 },
3359 quantity: Quantity::from("0.0001"),
3360 price: Some(Price::from("56730.0")),
3361 trigger_price: None,
3362 trigger_type: None,
3363 time_in_force: TimeInForce::Gtc,
3364 is_post_only: false,
3365 is_reduce_only: false,
3366 is_quote_quantity: false,
3367 }
3368 }
3369
3370 #[rstest]
3371 fn oid_query_attaches_the_known_client_order_id() {
3372 let mut report = make_status_report(None, "55030848197", OrderStatus::Accepted);
3373 let client_order_id = ClientOrderId::new("O-ATTACH-001");
3374
3375 attach_known_client_order_id(&mut report, client_order_id);
3376
3377 assert_eq!(report.client_order_id, Some(client_order_id));
3378 assert_eq!(report.venue_order_id, VenueOrderId::new("55030848197"));
3379 assert_eq!(report.order_status, OrderStatus::Accepted);
3380 }
3381
3382 #[rstest]
3383 fn oid_query_keeps_the_api_reported_client_order_id() {
3384 let mut report = make_status_report(
3385 Some("0x72a3c2f2de33c2c74640ad7f8d11ed74"),
3386 "222222",
3387 OrderStatus::Canceled,
3388 );
3389 let client_order_id = ClientOrderId::new("O-20240101-000002");
3390
3391 attach_known_client_order_id(&mut report, client_order_id);
3392
3393 assert_eq!(
3394 report.client_order_id,
3395 Some(ClientOrderId::new("0x72a3c2f2de33c2c74640ad7f8d11ed74"))
3396 );
3397 assert_eq!(report.venue_order_id, VenueOrderId::new("222222"));
3398 assert_eq!(report.order_status, OrderStatus::Canceled);
3399 }
3400
3401 fn make_status_report(
3402 client_order_id: Option<&str>,
3403 venue_order_id: &str,
3404 status: OrderStatus,
3405 ) -> OrderStatusReport {
3406 make_status_report_with_quantity(
3407 client_order_id,
3408 venue_order_id,
3409 status,
3410 Quantity::from("0.0001"),
3411 )
3412 }
3413
3414 fn make_status_report_with_quantity(
3415 client_order_id: Option<&str>,
3416 venue_order_id: &str,
3417 status: OrderStatus,
3418 quantity: Quantity,
3419 ) -> OrderStatusReport {
3420 OrderStatusReport::new(
3421 AccountId::from("HYPERLIQUID-001"),
3422 InstrumentId::from(TEST_INSTRUMENT_ID),
3423 client_order_id.map(ClientOrderId::new),
3424 VenueOrderId::new(venue_order_id),
3425 OrderSide::Buy,
3426 OrderType::Limit,
3427 TimeInForce::Gtc,
3428 status,
3429 quantity,
3430 Quantity::from("0"),
3431 UnixNanos::default(),
3432 UnixNanos::default(),
3433 UnixNanos::default(),
3434 Some(UUID4::new()),
3435 )
3436 .with_price(Price::from("56730.0"))
3437 }
3438
3439 fn make_fill_report(
3440 client_order_id: Option<&str>,
3441 venue_order_id: &str,
3442 trade_id: &str,
3443 ) -> FillReport {
3444 make_fill_report_with_qty(
3445 client_order_id,
3446 venue_order_id,
3447 trade_id,
3448 Quantity::from("0.0001"),
3449 )
3450 }
3451
3452 fn make_fill_report_with_qty(
3453 client_order_id: Option<&str>,
3454 venue_order_id: &str,
3455 trade_id: &str,
3456 last_qty: Quantity,
3457 ) -> FillReport {
3458 FillReport::new(
3459 AccountId::from("HYPERLIQUID-001"),
3460 InstrumentId::from(TEST_INSTRUMENT_ID),
3461 VenueOrderId::new(venue_order_id),
3462 TradeId::new(trade_id),
3463 OrderSide::Buy,
3464 last_qty,
3465 Price::from("56730.0"),
3466 Money::new(0.0, Currency::USD()),
3467 LiquiditySide::Taker,
3468 client_order_id.map(ClientOrderId::new),
3469 None,
3470 UnixNanos::default(),
3471 UnixNanos::default(),
3472 Some(UUID4::new()),
3473 )
3474 }
3475
3476 fn cloid_for(id: &str) -> Ustr {
3477 let cloid = Cloid::from_client_order_id(ClientOrderId::from(id));
3478 Ustr::from(&cloid.to_hex())
3479 }
3480
3481 #[rstest]
3482 fn test_filter_order_status_reports_for_command_filters_open_only() {
3483 let open_report =
3484 make_status_report(Some("O-HER-FILTER-OPEN"), "v-open", OrderStatus::Accepted);
3485 let closed_report =
3486 make_status_report(Some("O-HER-FILTER-CLOSED"), "v-closed", OrderStatus::Filled);
3487 let cmd = order_reports_command(true, None, None);
3488
3489 let filtered =
3490 filter_order_status_reports_for_command(vec![open_report, closed_report], &cmd);
3491
3492 assert_eq!(filtered.len(), 1);
3493 assert_eq!(
3494 filtered[0].client_order_id,
3495 Some(ClientOrderId::from("O-HER-FILTER-OPEN"))
3496 );
3497 }
3498
3499 #[rstest]
3500 fn test_filter_order_status_reports_for_command_filters_time_range_inclusively() {
3501 let mut before = make_status_report(
3502 Some("O-HER-FILTER-BEFORE"),
3503 "v-before",
3504 OrderStatus::Accepted,
3505 );
3506 let mut at_start =
3507 make_status_report(Some("O-HER-FILTER-START"), "v-start", OrderStatus::Accepted);
3508 let mut at_end =
3509 make_status_report(Some("O-HER-FILTER-END"), "v-end", OrderStatus::Accepted);
3510 let mut after =
3511 make_status_report(Some("O-HER-FILTER-AFTER"), "v-after", OrderStatus::Accepted);
3512 before.ts_last = UnixNanos::from(9);
3513 at_start.ts_last = UnixNanos::from(10);
3514 at_end.ts_last = UnixNanos::from(20);
3515 after.ts_last = UnixNanos::from(21);
3516 let cmd =
3517 order_reports_command(false, Some(UnixNanos::from(10)), Some(UnixNanos::from(20)));
3518
3519 let filtered =
3520 filter_order_status_reports_for_command(vec![before, at_start, at_end, after], &cmd);
3521 let filtered_ids: Vec<Option<ClientOrderId>> = filtered
3522 .iter()
3523 .map(|report| report.client_order_id)
3524 .collect();
3525
3526 assert_eq!(
3527 filtered_ids,
3528 vec![
3529 Some(ClientOrderId::from("O-HER-FILTER-START")),
3530 Some(ClientOrderId::from("O-HER-FILTER-END")),
3531 ]
3532 );
3533 }
3534
3535 #[rstest]
3536 fn test_filter_order_status_reports_for_command_without_filters_preserves_reports() {
3537 let open_report = make_status_report(
3538 Some("O-HER-FILTER-KEEP-OPEN"),
3539 "v-keep-open",
3540 OrderStatus::Accepted,
3541 );
3542 let closed_report = make_status_report(
3543 Some("O-HER-FILTER-KEEP-CLOSED"),
3544 "v-keep-closed",
3545 OrderStatus::Canceled,
3546 );
3547 let cmd = order_reports_command(false, None, None);
3548
3549 let filtered =
3550 filter_order_status_reports_for_command(vec![open_report, closed_report], &cmd);
3551 let filtered_ids: Vec<Option<ClientOrderId>> = filtered
3552 .iter()
3553 .map(|report| report.client_order_id)
3554 .collect();
3555
3556 assert_eq!(
3557 filtered_ids,
3558 vec![
3559 Some(ClientOrderId::from("O-HER-FILTER-KEEP-OPEN")),
3560 Some(ClientOrderId::from("O-HER-FILTER-KEEP-CLOSED")),
3561 ]
3562 );
3563 }
3564
3565 fn order_reports_command(
3566 open_only: bool,
3567 start: Option<UnixNanos>,
3568 end: Option<UnixNanos>,
3569 ) -> GenerateOrderStatusReports {
3570 GenerateOrderStatusReports::new(
3571 UUID4::new(),
3572 UnixNanos::default(),
3573 open_only,
3574 None,
3575 start,
3576 end,
3577 None,
3578 None,
3579 )
3580 }
3581
3582 fn limit_order(
3583 id: &str,
3584 reduce_only: bool,
3585 contingency: ContingencyType,
3586 linked_ids: Option<Vec<&str>>,
3587 parent_id: Option<&str>,
3588 ) -> OrderAny {
3589 OrderAny::Limit(LimitOrder::new(
3590 TraderId::from("TESTER-001"),
3591 StrategyId::from("S-001"),
3592 InstrumentId::from("ETH-USD-PERP.HYPERLIQUID"),
3593 ClientOrderId::from(id),
3594 OrderSide::Buy,
3595 Quantity::from(1),
3596 Price::from("3000.00"),
3597 TimeInForce::Gtc,
3598 None, false, reduce_only,
3601 false, None, None, None, Some(contingency),
3606 None, linked_ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()),
3608 parent_id.map(ClientOrderId::from),
3609 None, None, None, None, Default::default(),
3614 Default::default(),
3615 ))
3616 }
3617
3618 fn stop_order(
3619 id: &str,
3620 reduce_only: bool,
3621 contingency: ContingencyType,
3622 linked_ids: Option<Vec<&str>>,
3623 parent_id: Option<&str>,
3624 ) -> OrderAny {
3625 OrderAny::StopMarket(StopMarketOrder::new(
3626 TraderId::from("TESTER-001"),
3627 StrategyId::from("S-001"),
3628 InstrumentId::from("ETH-USD-PERP.HYPERLIQUID"),
3629 ClientOrderId::from(id),
3630 OrderSide::Sell,
3631 Quantity::from(1),
3632 Price::from("2800.00"),
3633 TriggerType::LastPrice,
3634 TimeInForce::Gtc,
3635 None, reduce_only,
3637 false, None, None, None, Some(contingency),
3642 None, linked_ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()),
3644 parent_id.map(ClientOrderId::from),
3645 None, None, None, None, Default::default(),
3650 Default::default(),
3651 ))
3652 }
3653
3654 fn staged_child(id: &str, sibling_id: &str) -> StagedBracketChild {
3655 StagedBracketChild {
3656 order: limit_order(
3657 id,
3658 true,
3659 ContingencyType::Ouo,
3660 Some(vec![sibling_id]),
3661 Some("O-PARENT"),
3662 ),
3663 request: HyperliquidExecPlaceOrderRequest {
3664 asset: 4,
3665 is_buy: false,
3666 price: Decimal::from(3_000),
3667 size: Decimal::ONE,
3668 reduce_only: true,
3669 kind: HyperliquidExecOrderKind::Limit {
3670 limit: HyperliquidExecLimitParams {
3671 tif: HyperliquidExecTif::Gtc,
3672 },
3673 },
3674 cloid: Some(Cloid::from_client_order_id(ClientOrderId::from(id))),
3675 },
3676 }
3677 }
3678
3679 #[rstest]
3680 fn test_staged_bracket_activation_links_ouo_siblings_once() {
3681 let parent_id = ClientOrderId::from("O-PARENT");
3682 let first_id = ClientOrderId::from("O-CHILD-1");
3683 let second_id = ClientOrderId::from("O-CHILD-2");
3684 let mut state = StagedBracketState::default();
3685 state.stage(
3686 parent_id,
3687 vec![
3688 staged_child(first_id.as_str(), second_id.as_str()),
3689 staged_child(second_id.as_str(), first_id.as_str()),
3690 ],
3691 );
3692
3693 let activated = state.activate(&parent_id).expect("staged children");
3694 let sibling = state
3695 .take_active_sibling(&first_id)
3696 .expect("active OUO sibling");
3697
3698 assert_eq!(activated.len(), 2);
3699 assert_eq!(sibling.order.client_order_id(), second_id);
3700 assert!(state.activate(&parent_id).is_none());
3701 assert!(state.take_active_sibling(&second_id).is_none());
3702 }
3703
3704 #[rstest]
3705 fn test_restored_active_bracket_rebuilds_ouo_without_reactivation() {
3706 let parent_id = ClientOrderId::from("O-PARENT");
3707 let first_id = ClientOrderId::from("O-CHILD-1");
3708 let second_id = ClientOrderId::from("O-CHILD-2");
3709 let mut state = StagedBracketState::default();
3710 state.restore_active(&[
3711 staged_child(first_id.as_str(), second_id.as_str()),
3712 staged_child(second_id.as_str(), first_id.as_str()),
3713 ]);
3714
3715 let sibling = state
3716 .take_active_sibling(&first_id)
3717 .expect("restored OUO sibling");
3718
3719 assert!(state.activate(&parent_id).is_none());
3720 assert_eq!(sibling.order.client_order_id(), second_id);
3721 assert!(state.take_active_sibling(&second_id).is_none());
3722 }
3723
3724 #[rstest]
3725 fn test_staged_bracket_child_cancel_preserves_other_child_for_parent_fill() {
3726 let parent_id = ClientOrderId::from("O-PARENT");
3727 let first_id = ClientOrderId::from("O-CHILD-1");
3728 let second_id = ClientOrderId::from("O-CHILD-2");
3729 let mut state = StagedBracketState::default();
3730 state.stage(
3731 parent_id,
3732 vec![
3733 staged_child(first_id.as_str(), second_id.as_str()),
3734 staged_child(second_id.as_str(), first_id.as_str()),
3735 ],
3736 );
3737
3738 let canceled = state.cancel_child(&first_id).expect("staged child");
3739 let remaining = state.activate(&parent_id).expect("remaining child");
3740
3741 assert_eq!(canceled.client_order_id(), first_id);
3742 assert_eq!(remaining.len(), 1);
3743 assert_eq!(remaining[0].order.client_order_id(), second_id);
3744 }
3745
3746 #[rstest]
3747 fn test_staged_bracket_parent_cancel_returns_all_unsubmitted_children() {
3748 let parent_id = ClientOrderId::from("O-PARENT");
3749 let first_id = ClientOrderId::from("O-CHILD-1");
3750 let second_id = ClientOrderId::from("O-CHILD-2");
3751 let mut state = StagedBracketState::default();
3752 state.stage(
3753 parent_id,
3754 vec![
3755 staged_child(first_id.as_str(), second_id.as_str()),
3756 staged_child(second_id.as_str(), first_id.as_str()),
3757 ],
3758 );
3759
3760 let canceled = state.cancel_for_parent(&parent_id);
3761 let canceled_ids = canceled
3762 .iter()
3763 .map(Order::client_order_id)
3764 .collect::<Vec<_>>();
3765
3766 assert_eq!(canceled_ids, vec![first_id, second_id]);
3767 assert!(state.activate(&parent_id).is_none());
3768 }
3769
3770 #[rstest]
3771 fn test_build_ouo_resize_request_sends_sibling_leaves_quantity() {
3772 let sibling = staged_child("O-CHILD-2", "O-CHILD-1");
3773
3774 let request =
3775 build_ouo_resize_request(&sibling, Quantity::from("0.7"), Quantity::from("0.2"))
3776 .expect("resized request");
3777 let exhausted =
3778 build_ouo_resize_request(&sibling, Quantity::from("0.2"), Quantity::from("0.2"));
3779
3780 assert_eq!(request.size, Decimal::new(5, 1));
3781 assert_eq!(request.cloid, sibling.request.cloid);
3782 assert!(exhausted.is_none());
3783 }
3784
3785 #[rstest]
3786 #[case::independent_orders(
3787 vec![
3788 limit_order("O-001", false, ContingencyType::NoContingency, None, None),
3789 limit_order("O-002", false, ContingencyType::NoContingency, None, None),
3790 ],
3791 HyperliquidExecGrouping::Na,
3792 )]
3793 #[case::bracket_oto(
3794 vec![
3795 limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002", "O-003"]), None),
3796 limit_order("O-002", true, ContingencyType::Oco, Some(vec!["O-003"]), Some("O-001")),
3797 stop_order("O-003", true, ContingencyType::Oco, Some(vec!["O-002"]), Some("O-001")),
3798 ],
3799 HyperliquidExecGrouping::NormalTpsl,
3800 )]
3801 #[case::bracket_oto_with_factory_ouo_children(
3802 vec![
3803 limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002", "O-003"]), None),
3804 limit_order("O-002", true, ContingencyType::Ouo, Some(vec!["O-003"]), Some("O-001")),
3805 stop_order("O-003", true, ContingencyType::Ouo, Some(vec!["O-002"]), Some("O-001")),
3806 ],
3807 HyperliquidExecGrouping::NormalTpsl,
3808 )]
3809 #[case::oto_not_bracket_shaped(
3810 vec![
3811 limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002"]), None),
3812 limit_order("O-002", false, ContingencyType::Oto, Some(vec!["O-001"]), None),
3813 ],
3814 HyperliquidExecGrouping::Na,
3815 )]
3816 #[case::oco_all_reduce_only(
3817 vec![
3818 limit_order("O-001", true, ContingencyType::Oco, Some(vec!["O-002"]), None),
3819 stop_order("O-002", true, ContingencyType::Oco, Some(vec!["O-001"]), None),
3820 ],
3821 HyperliquidExecGrouping::PositionTpsl,
3822 )]
3823 #[case::oco_not_all_reduce_only(
3824 vec![
3825 limit_order("O-001", false, ContingencyType::Oco, Some(vec!["O-002"]), None),
3826 stop_order("O-002", true, ContingencyType::Oco, Some(vec!["O-001"]), None),
3827 ],
3828 HyperliquidExecGrouping::Na,
3829 )]
3830 #[case::oto_with_non_oco_children(
3831 vec![
3832 limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002", "O-003"]), None),
3833 limit_order("O-002", true, ContingencyType::NoContingency, None, None),
3834 stop_order("O-003", true, ContingencyType::NoContingency, None, None),
3835 ],
3836 HyperliquidExecGrouping::Na,
3837 )]
3838 #[case::mixed_oco_and_plain_reduce_only(
3839 vec![
3840 limit_order("O-001", true, ContingencyType::Oco, Some(vec!["O-002"]), None),
3841 stop_order("O-002", true, ContingencyType::NoContingency, None, None),
3842 ],
3843 HyperliquidExecGrouping::Na,
3844 )]
3845 #[case::unlinked_oco_reduce_only(
3846 vec![
3847 limit_order("O-001", true, ContingencyType::Oco, Some(vec!["O-099"]), None),
3848 stop_order("O-002", true, ContingencyType::Oco, Some(vec!["O-098"]), None),
3849 ],
3850 HyperliquidExecGrouping::Na,
3851 )]
3852 #[case::single_order(
3853 vec![limit_order("O-001", false, ContingencyType::NoContingency, None, None)],
3854 HyperliquidExecGrouping::Na,
3855 )]
3856 fn test_determine_order_list_grouping(
3857 #[case] orders: Vec<OrderAny>,
3858 #[case] expected: HyperliquidExecGrouping,
3859 ) {
3860 let result = determine_order_list_grouping(&orders);
3861 assert_eq!(result, expected);
3862 }
3863
3864 #[rstest]
3865 #[case::market(Some(OrderType::Market), true)]
3866 #[case::limit(Some(OrderType::Limit), true)]
3867 #[case::stop_market(Some(OrderType::StopMarket), false)]
3868 #[case::unknown(None, false)]
3869 fn test_can_fast_cancel_order_only_allows_plain_order_types(
3870 #[case] order_type: Option<OrderType>,
3871 #[case] expected: bool,
3872 ) {
3873 assert_eq!(can_fast_cancel_order(order_type), expected);
3874 }
3875
3876 #[rstest]
3877 fn test_split_fast_cancel_requests_preserves_request_entry_alignment() {
3878 let requests = vec![
3879 (10_u64, cancel_entry("O-FAST-1", true)),
3880 (20_u64, cancel_entry("O-NORMAL-1", false)),
3881 (30_u64, cancel_entry("O-FAST-2", true)),
3882 (40_u64, cancel_entry("O-NORMAL-2", false)),
3883 ];
3884
3885 let (fast_requests, fast_entries, normal_requests, normal_entries) =
3886 split_fast_cancel_requests(requests);
3887
3888 assert_eq!(fast_requests, vec![10, 30]);
3889 assert_eq!(
3890 client_order_ids(&fast_entries),
3891 vec![
3892 ClientOrderId::from("O-FAST-1"),
3893 ClientOrderId::from("O-FAST-2"),
3894 ]
3895 );
3896 assert!(fast_entries.iter().all(|entry| entry.fast));
3897 assert_eq!(normal_requests, vec![20, 40]);
3898 assert_eq!(
3899 client_order_ids(&normal_entries),
3900 vec![
3901 ClientOrderId::from("O-NORMAL-1"),
3902 ClientOrderId::from("O-NORMAL-2"),
3903 ]
3904 );
3905 assert!(normal_entries.iter().all(|entry| !entry.fast));
3906 }
3907
3908 fn cancel_entry(client_order_id: &str, fast: bool) -> CancelEntry {
3909 CancelEntry {
3910 strategy_id: StrategyId::from("S-001"),
3911 instrument_id: InstrumentId::from(TEST_INSTRUMENT_ID),
3912 client_order_id: ClientOrderId::from(client_order_id),
3913 venue_order_id: Some(VenueOrderId::new("123")),
3914 symbol: Ustr::from("BTC-USD-PERP"),
3915 fast,
3916 }
3917 }
3918
3919 fn client_order_ids(entries: &[CancelEntry]) -> Vec<ClientOrderId> {
3920 entries.iter().map(|entry| entry.client_order_id).collect()
3921 }
3922
3923 fn limit_order_with_flags(id: &str, quote_quantity: bool, post_only: bool) -> OrderAny {
3924 OrderAny::Limit(LimitOrder::new(
3925 TraderId::from("TESTER-001"),
3926 StrategyId::from("S-001"),
3927 InstrumentId::from(TEST_INSTRUMENT_ID),
3928 ClientOrderId::from(id),
3929 OrderSide::Buy,
3930 Quantity::from("0.0001"),
3931 Price::from("56730.0"),
3932 TimeInForce::Gtc,
3933 None,
3934 post_only,
3935 false,
3936 quote_quantity,
3937 None,
3938 None,
3939 None,
3940 Some(ContingencyType::NoContingency),
3941 None,
3942 None,
3943 None,
3944 None,
3945 None,
3946 None,
3947 None,
3948 Default::default(),
3949 Default::default(),
3950 ))
3951 }
3952
3953 #[rstest]
3954 fn test_register_order_context_registers_regular_order() {
3955 let state = WsDispatchState::new();
3956 let client_order_id = ClientOrderId::from("O-REG-001");
3957 let order = limit_order_with_flags("O-REG-001", false, false);
3958
3959 register_order_context_into(&state, &order);
3960
3961 assert_eq!(
3962 state.lookup_context(&client_order_id),
3963 Some(test_context(client_order_id)),
3964 );
3965 }
3966
3967 #[rstest]
3968 fn test_register_order_context_skips_quote_quantity_order() {
3969 let state = WsDispatchState::new();
3970 let order = limit_order_with_flags("O-QQ-001", true, false);
3971
3972 register_order_context_into(&state, &order);
3973
3974 assert!(
3979 state
3980 .lookup_context(&ClientOrderId::from("O-QQ-001"))
3981 .is_none()
3982 );
3983 }
3984
3985 #[rstest]
3986 fn test_handle_execution_report_skip_keeps_cloid_mapping() {
3987 let ws_client = make_ws_client();
3992 let (emitter, mut rx) = test_emitter();
3993 let state = WsDispatchState::new();
3994 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
3995
3996 let cid = ClientOrderId::from("O-HER-SKIP");
3997 state.register_context(test_context(cid));
3998 state.insert_accepted(cid);
4000 state.record_venue_order_id(cid, VenueOrderId::new("new-voi"));
4001
4002 ws_client.cache_cloid_mapping(cloid_for("O-HER-SKIP"), cid);
4003
4004 let stale_cancel = make_status_report(Some("O-HER-SKIP"), "old-voi", OrderStatus::Canceled);
4005 handle_execution_report(
4006 ExecutionReport::Order(stale_cancel),
4007 &state,
4008 &emitter,
4009 &ws_client,
4010 &make_http_client(),
4011 &mut pending_cloids,
4012 UnixNanos::default(),
4013 );
4014
4015 assert!(drain_events(&mut rx).is_empty());
4016 assert_eq!(
4018 ws_client.get_cloid_mapping(&cloid_for("O-HER-SKIP")),
4019 Some(cid)
4020 );
4021 assert!(state.lookup_context(&cid).is_some());
4023 }
4024
4025 #[rstest]
4026 fn test_handle_execution_report_tracked_terminal_evicts_cloid() {
4027 let ws_client = make_ws_client();
4031 let (emitter, mut rx) = test_emitter();
4032 let state = WsDispatchState::new();
4033 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4034
4035 let cid = ClientOrderId::from("O-HER-CANCEL");
4036 state.register_context(test_context(cid));
4037 state.insert_accepted(cid);
4038 state.record_venue_order_id(cid, VenueOrderId::new("v-cancel"));
4039
4040 ws_client.cache_cloid_mapping(cloid_for("O-HER-CANCEL"), cid);
4041
4042 let report = make_status_report(Some("O-HER-CANCEL"), "v-cancel", OrderStatus::Canceled);
4043 handle_execution_report(
4044 ExecutionReport::Order(report),
4045 &state,
4046 &emitter,
4047 &ws_client,
4048 &make_http_client(),
4049 &mut pending_cloids,
4050 UnixNanos::default(),
4051 );
4052
4053 let events = drain_events(&mut rx);
4054 assert_eq!(events.len(), 1);
4055 assert!(matches!(
4056 events[0],
4057 ExecutionEvent::Order(OrderEventAny::Canceled(_))
4058 ));
4059 assert_eq!(
4060 ws_client.get_cloid_mapping(&cloid_for("O-HER-CANCEL")),
4061 None
4062 );
4063 assert!(state.filled_orders.contains(&cid));
4064 }
4065
4066 #[rstest]
4067 fn test_post_rejection_preserves_exact_reason_when_ws_rejection_arrives_first() {
4068 let ws_client = make_ws_client();
4069 let (emitter, mut rx) = test_emitter();
4070 let state = Arc::new(WsDispatchState::new());
4071 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4072
4073 let cid = ClientOrderId::from("O-HER-WS-REJ");
4074 state.register_context(test_context(cid));
4075 state.mark_submission_pending(cid);
4076 ws_client.cache_cloid_mapping(cloid_for("O-HER-WS-REJ"), cid);
4077
4078 let report = make_status_report(Some("O-HER-WS-REJ"), "v-rej", OrderStatus::Rejected);
4079 handle_execution_report(
4080 ExecutionReport::Order(report),
4081 &state,
4082 &emitter,
4083 &ws_client,
4084 &make_http_client(),
4085 &mut pending_cloids,
4086 UnixNanos::default(),
4087 );
4088
4089 assert!(drain_events(&mut rx).is_empty());
4090 assert_eq!(
4091 ws_client.get_cloid_mapping(&cloid_for("O-HER-WS-REJ")),
4092 Some(cid),
4093 );
4094
4095 let order = limit_order_with_flags("O-HER-WS-REJ", false, true);
4096 let http_client = make_http_client();
4097 let rejection_route =
4098 PostRejectionRoute::new(&emitter, &ws_client, &http_client, state.clone());
4099 let emitted = rejection_route.emit_once(
4100 &order,
4101 "Post only order would have immediately matched, bbo was 56729.0.",
4102 UnixNanos::default(),
4103 &cloid_for("O-HER-WS-REJ"),
4104 );
4105
4106 let events = drain_events(&mut rx);
4107 let ExecutionEvent::Order(OrderEventAny::Rejected(rejected)) = &events[0] else {
4108 panic!("expected OrderRejected, received {:?}", events[0]);
4109 };
4110 assert!(emitted);
4111 assert_eq!(events.len(), 1);
4112 assert_eq!(
4113 rejected.reason.as_str(),
4114 "Post only order would have immediately matched, bbo was 56729.0.",
4115 );
4116 assert!(rejected.due_post_only);
4117 }
4118
4119 #[rstest]
4120 fn test_post_rejection_suppresses_late_raw_cloid_reject() {
4121 let ws_client = make_ws_client();
4122 let (emitter, mut rx) = test_emitter();
4123 let state = Arc::new(WsDispatchState::new());
4124 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4125
4126 let cid = ClientOrderId::from("O-HER-POST-REJ");
4127 let cloid = cloid_for("O-HER-POST-REJ");
4128 let order = limit_order_with_flags("O-HER-POST-REJ", false, true);
4129 state.register_context(test_context(cid));
4130 ws_client.cache_cloid_mapping(cloid, cid);
4131
4132 let http_client = make_http_client();
4133 let rejection_route =
4134 PostRejectionRoute::new(&emitter, &ws_client, &http_client, state.clone());
4135 let emitted = rejection_route.emit_once(
4136 &order,
4137 "Post only order would have immediately matched",
4138 UnixNanos::default(),
4139 &cloid,
4140 );
4141
4142 let events = drain_events(&mut rx);
4143 assert!(emitted);
4144 assert_eq!(events.len(), 1);
4145 let ExecutionEvent::Order(OrderEventAny::Rejected(rejected)) = &events[0] else {
4146 panic!("expected OrderRejected, received {:?}", events[0]);
4147 };
4148 assert_eq!(
4149 rejected.reason.as_str(),
4150 "Post only order would have immediately matched",
4151 );
4152 assert!(rejected.due_post_only);
4153 assert_eq!(ws_client.get_cloid_mapping(&cloid), None);
4154 assert!(state.filled_orders.contains(&cid));
4155 assert!(state.terminal_cloid_seen(&cloid));
4156
4157 let late_reject = make_status_report(Some(cloid.as_str()), "v-rej", OrderStatus::Rejected);
4158 handle_execution_report(
4159 ExecutionReport::Order(late_reject),
4160 &state,
4161 &emitter,
4162 &ws_client,
4163 &make_http_client(),
4164 &mut pending_cloids,
4165 UnixNanos::default(),
4166 );
4167
4168 assert!(drain_events(&mut rx).is_empty());
4169 }
4170
4171 #[rstest]
4172 fn test_handle_execution_report_filled_marker_then_fill_evicts_on_fill() {
4173 let ws_client = make_ws_client();
4177 let (emitter, mut rx) = test_emitter();
4178 let state = WsDispatchState::new();
4179 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4180
4181 let cid = ClientOrderId::from("O-HER-FILL");
4182 state.register_context(test_context(cid));
4183 state.insert_accepted(cid);
4184 state.record_venue_order_id(cid, VenueOrderId::new("v-fill"));
4185
4186 ws_client.cache_cloid_mapping(cloid_for("O-HER-FILL"), cid);
4187
4188 let status_marker = make_status_report(Some("O-HER-FILL"), "v-fill", OrderStatus::Filled);
4189 handle_execution_report(
4190 ExecutionReport::Order(status_marker),
4191 &state,
4192 &emitter,
4193 &ws_client,
4194 &make_http_client(),
4195 &mut pending_cloids,
4196 UnixNanos::default(),
4197 );
4198
4199 assert!(drain_events(&mut rx).is_empty());
4201 assert_eq!(
4202 ws_client.get_cloid_mapping(&cloid_for("O-HER-FILL")),
4203 Some(cid)
4204 );
4205
4206 let fill = make_fill_report(Some("O-HER-FILL"), "v-fill", "trade-fill");
4207 handle_execution_report(
4208 ExecutionReport::Fill(fill),
4209 &state,
4210 &emitter,
4211 &ws_client,
4212 &make_http_client(),
4213 &mut pending_cloids,
4214 UnixNanos::default(),
4215 );
4216
4217 let events = drain_events(&mut rx);
4218 assert_eq!(events.len(), 1);
4219 assert!(matches!(
4220 events[0],
4221 ExecutionEvent::Order(OrderEventAny::Filled(_))
4222 ));
4223 assert_eq!(ws_client.get_cloid_mapping(&cloid_for("O-HER-FILL")), None);
4225 }
4226
4227 #[rstest]
4232 fn test_handle_execution_report_fill_under_filled_marker_promotes_and_evicts_cloid() {
4233 let ws_client = make_ws_client();
4234 let (emitter, mut rx) = test_emitter();
4235 let state = WsDispatchState::new();
4236 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4237
4238 let cid = ClientOrderId::from("O-HER-BUF");
4239 state.register_context(test_context(cid));
4240 state.insert_accepted(cid);
4241 state.record_venue_order_id(cid, VenueOrderId::new("old-voi"));
4242 state.mark_pending_modify(
4243 cid,
4244 VenueOrderId::new("old-voi"),
4245 test_context(cid).quantity,
4246 );
4247
4248 ws_client.cache_cloid_mapping(cloid_for("O-HER-BUF"), cid);
4249
4250 let status_marker = make_status_report(Some("O-HER-BUF"), "new-voi", OrderStatus::Filled);
4252 handle_execution_report(
4253 ExecutionReport::Order(status_marker),
4254 &state,
4255 &emitter,
4256 &ws_client,
4257 &make_http_client(),
4258 &mut pending_cloids,
4259 UnixNanos::default(),
4260 );
4261 assert!(pending_cloids.contains(&cid));
4262 assert_eq!(
4263 ws_client.get_cloid_mapping(&cloid_for("O-HER-BUF")),
4264 Some(cid)
4265 );
4266
4267 let fill = make_fill_report(Some("O-HER-BUF"), "new-voi", "trade-buf");
4271 handle_execution_report(
4272 ExecutionReport::Fill(fill),
4273 &state,
4274 &emitter,
4275 &ws_client,
4276 &make_http_client(),
4277 &mut pending_cloids,
4278 UnixNanos::default(),
4279 );
4280
4281 let events = drain_events(&mut rx);
4282 assert_eq!(events.len(), 2);
4283 assert!(matches!(
4284 events[0],
4285 ExecutionEvent::Order(OrderEventAny::Updated(_))
4286 ));
4287 assert!(matches!(
4288 events[1],
4289 ExecutionEvent::Order(OrderEventAny::Filled(_))
4290 ));
4291 assert_eq!(state.buffered_fill_count(&cid), 0);
4292 assert!(
4293 !pending_cloids.contains(&cid),
4294 "deferred cleanup must complete once the promoting fill lands",
4295 );
4296 assert_eq!(
4297 ws_client.get_cloid_mapping(&cloid_for("O-HER-BUF")),
4298 None,
4299 "cloid mapping must be evicted after the terminal fill",
4300 );
4301 }
4302
4303 #[rstest]
4306 fn test_cancel_replace_emits_target_total_quantity() {
4307 let ws_client = make_ws_client();
4308 let (emitter, mut rx) = test_emitter();
4309 let state = WsDispatchState::new();
4310 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4311
4312 let cid = ClientOrderId::from("O-HER-CR-QTY");
4313 let target_total = Quantity::from("0.00020");
4314 let venue_remaining = Quantity::from("0.00015");
4315
4316 let mut context = test_context(cid);
4317 context.quantity = target_total;
4318 state.register_context(context);
4319 state.insert_accepted(cid);
4320 state.record_venue_order_id(cid, VenueOrderId::new("old-voi"));
4321 state.mark_pending_modify(cid, VenueOrderId::new("old-voi"), target_total);
4322
4323 ws_client.cache_cloid_mapping(cloid_for("O-HER-CR-QTY"), cid);
4324
4325 let accepted = make_status_report_with_quantity(
4326 Some("O-HER-CR-QTY"),
4327 "new-voi",
4328 OrderStatus::Accepted,
4329 venue_remaining,
4330 );
4331 handle_execution_report(
4332 ExecutionReport::Order(accepted),
4333 &state,
4334 &emitter,
4335 &ws_client,
4336 &make_http_client(),
4337 &mut pending_cloids,
4338 UnixNanos::default(),
4339 );
4340
4341 let events = drain_events(&mut rx);
4342 assert_eq!(events.len(), 1);
4343 match &events[0] {
4344 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
4345 assert_eq!(
4346 updated.quantity, target_total,
4347 "OrderUpdated must carry the engine's absolute total quantity",
4348 );
4349 assert_eq!(updated.venue_order_id, Some(VenueOrderId::new("new-voi")));
4350 }
4351 other => panic!("expected OrderUpdated, found {other:?}"),
4352 }
4353
4354 let context = state
4356 .lookup_context(&cid)
4357 .expect("context should still be tracked");
4358 assert_eq!(context.quantity, target_total);
4359
4360 assert!(state.pending_modify(&cid).is_none());
4361 assert!(state.pending_modify_target_qty(&cid).is_none());
4362 assert_eq!(
4363 state.cached_venue_order_id(&cid),
4364 Some(VenueOrderId::new("new-voi")),
4365 );
4366 }
4367
4368 #[rstest]
4371 fn test_cancel_replace_without_marker_falls_back_to_report_quantity() {
4372 let ws_client = make_ws_client();
4373 let (emitter, mut rx) = test_emitter();
4374 let state = WsDispatchState::new();
4375 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4376
4377 let cid = ClientOrderId::from("O-HER-CR-EXT");
4378 state.register_context(test_context(cid));
4379 state.insert_accepted(cid);
4380 state.record_venue_order_id(cid, VenueOrderId::new("old-voi"));
4381
4382 ws_client.cache_cloid_mapping(cloid_for("O-HER-CR-EXT"), cid);
4383
4384 let report_qty = Quantity::from("0.0005");
4385 let accepted = make_status_report_with_quantity(
4386 Some("O-HER-CR-EXT"),
4387 "new-voi",
4388 OrderStatus::Accepted,
4389 report_qty,
4390 );
4391 handle_execution_report(
4392 ExecutionReport::Order(accepted),
4393 &state,
4394 &emitter,
4395 &ws_client,
4396 &make_http_client(),
4397 &mut pending_cloids,
4398 UnixNanos::default(),
4399 );
4400
4401 let events = drain_events(&mut rx);
4402 assert_eq!(events.len(), 1);
4403 match &events[0] {
4404 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
4405 assert_eq!(updated.quantity, report_qty);
4406 }
4407 other => panic!("expected OrderUpdated, found {other:?}"),
4408 }
4409 }
4410
4411 fn limit_request(size: Decimal) -> HyperliquidExecPlaceOrderRequest {
4412 HyperliquidExecPlaceOrderRequest {
4413 asset: 0,
4414 is_buy: true,
4415 price: "88.949".parse::<Decimal>().unwrap(),
4416 size,
4417 reduce_only: false,
4418 kind: HyperliquidExecOrderKind::Limit {
4419 limit: HyperliquidExecLimitParams {
4420 tif: HyperliquidExecTif::Gtc,
4421 },
4422 },
4423 cloid: None,
4424 }
4425 }
4426
4427 #[rstest]
4431 fn test_cancel_replace_queues_corrective_reduce_on_in_flight_fill() {
4432 let ws_client = make_ws_client();
4433 let (emitter, mut rx) = test_emitter();
4434 let state = WsDispatchState::new();
4435 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4436
4437 let cid = ClientOrderId::from("O-HER-4154");
4438 let target_total = Quantity::from("1.000");
4439 let old_voi = "445117664938";
4440 let new_voi = "445117686214";
4441
4442 let mut context = test_context(cid);
4443 context.quantity = target_total;
4444 state.register_context(context);
4445 state.insert_accepted(cid);
4446 state.record_venue_order_id(cid, VenueOrderId::new(old_voi));
4447
4448 state.mark_pending_modify(cid, VenueOrderId::new(old_voi), target_total);
4451 state.stash_modify_request(cid, limit_request(Decimal::from(1)));
4452
4453 state.record_filled_qty(cid, Quantity::from("0.165"));
4455
4456 let accepted = make_status_report_with_quantity(
4458 Some("O-HER-4154"),
4459 new_voi,
4460 OrderStatus::Accepted,
4461 Quantity::from("0.835"),
4462 );
4463 let corrective = handle_execution_report(
4464 ExecutionReport::Order(accepted),
4465 &state,
4466 &emitter,
4467 &ws_client,
4468 &make_http_client(),
4469 &mut pending_cloids,
4470 UnixNanos::default(),
4471 );
4472
4473 let events = drain_events(&mut rx);
4475 assert_eq!(events.len(), 1);
4476 match &events[0] {
4477 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
4478 assert_eq!(updated.quantity, target_total);
4479 assert_eq!(updated.venue_order_id, Some(VenueOrderId::new(new_voi)));
4480 }
4481 other => panic!("expected OrderUpdated, found {other:?}"),
4482 }
4483
4484 let (corr_cid, oid, request) =
4485 corrective.expect("oversized replacement must queue a corrective reduce");
4486 assert_eq!(corr_cid, cid);
4487 assert_eq!(oid, 445_117_686_214);
4488 assert_eq!(request.size, "0.835".parse::<Decimal>().unwrap());
4489 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new(new_voi)));
4492 assert_eq!(state.pending_modify_target_qty(&cid), Some(target_total));
4493 }
4494
4495 #[rstest]
4499 fn test_cancel_replace_fill_promotion_queues_corrective_reduce() {
4500 let ws_client = make_ws_client();
4501 let (emitter, mut rx) = test_emitter();
4502 let state = WsDispatchState::new();
4503 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4504
4505 let cid = ClientOrderId::from("O-HER-FILL-CORR");
4506 let target_total = Quantity::from("1.000");
4507 let old_voi = "445117664938";
4508 let new_voi = "445117686214";
4509
4510 let mut context = test_context(cid);
4511 context.quantity = target_total;
4512 state.register_context(context);
4513 state.insert_accepted(cid);
4514 state.record_venue_order_id(cid, VenueOrderId::new(old_voi));
4515 state.mark_pending_modify(cid, VenueOrderId::new(old_voi), target_total);
4517 state.stash_modify_request(cid, limit_request(Decimal::from(1)));
4518 state.record_filled_qty(cid, Quantity::from("0.165"));
4520
4521 let fill = make_fill_report_with_qty(
4523 Some("O-HER-FILL-CORR"),
4524 new_voi,
4525 "T-FILL-CORR",
4526 Quantity::from("0.100"),
4527 );
4528 let corrective = handle_execution_report(
4529 ExecutionReport::Fill(fill),
4530 &state,
4531 &emitter,
4532 &ws_client,
4533 &make_http_client(),
4534 &mut pending_cloids,
4535 UnixNanos::default(),
4536 );
4537
4538 let events = drain_events(&mut rx);
4540 assert_eq!(events.len(), 2);
4541 assert!(matches!(
4542 events[0],
4543 ExecutionEvent::Order(OrderEventAny::Updated(_))
4544 ));
4545 assert!(matches!(
4546 events[1],
4547 ExecutionEvent::Order(OrderEventAny::Filled(_))
4548 ));
4549
4550 let (corr_cid, oid, request) =
4552 corrective.expect("oversized replacement must queue a corrective reduce");
4553 assert_eq!(corr_cid, cid);
4554 assert_eq!(oid, 445_117_686_214);
4555 assert_eq!(request.size, "0.735".parse::<Decimal>().unwrap());
4556 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new(new_voi)));
4557 }
4558
4559 #[rstest]
4562 fn test_cancel_replace_no_corrective_without_in_flight_fill() {
4563 let ws_client = make_ws_client();
4564 let (emitter, mut rx) = test_emitter();
4565 let state = WsDispatchState::new();
4566 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4567
4568 let cid = ClientOrderId::from("O-HER-4154-NOFILL");
4569 let target_total = Quantity::from("1.000");
4570
4571 let mut context = test_context(cid);
4572 context.quantity = target_total;
4573 state.register_context(context);
4574 state.insert_accepted(cid);
4575 state.record_venue_order_id(cid, VenueOrderId::new("445117664938"));
4576 state.mark_pending_modify(cid, VenueOrderId::new("445117664938"), target_total);
4577 state.stash_modify_request(cid, limit_request(Decimal::from(1)));
4578
4579 let accepted = make_status_report_with_quantity(
4580 Some("O-HER-4154-NOFILL"),
4581 "445117686214",
4582 OrderStatus::Accepted,
4583 target_total,
4584 );
4585 let corrective = handle_execution_report(
4586 ExecutionReport::Order(accepted),
4587 &state,
4588 &emitter,
4589 &ws_client,
4590 &make_http_client(),
4591 &mut pending_cloids,
4592 UnixNanos::default(),
4593 );
4594
4595 let _ = drain_events(&mut rx);
4596 assert!(corrective.is_none());
4597 assert!(state.pending_modify(&cid).is_none());
4598 assert!(state.take_corrective(&cid).is_none());
4599 assert!(state.modify_request(&cid).is_none());
4601 }
4602
4603 #[rstest]
4607 fn test_cancel_replace_corrective_uses_post_drain_buffered_fill() {
4608 let ws_client = make_ws_client();
4609 let (emitter, mut rx) = test_emitter();
4610 let state = WsDispatchState::new();
4611 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4612
4613 let cid = ClientOrderId::from("O-HER-4154-BUF");
4614 let target_total = Quantity::from("1.000");
4615 let new_voi = "445117686214";
4616
4617 let mut context = test_context(cid);
4618 context.quantity = target_total;
4619 state.register_context(context);
4620 state.insert_accepted(cid);
4621 state.record_venue_order_id(cid, VenueOrderId::new("445117664938"));
4622 state.mark_pending_modify(cid, VenueOrderId::new("445117664938"), target_total);
4623 state.stash_modify_request(cid, limit_request(Decimal::from(1)));
4624
4625 let buffered = make_fill_report_with_qty(
4627 Some("O-HER-4154-BUF"),
4628 new_voi,
4629 "trade-buf-4154",
4630 Quantity::from("0.165"),
4631 );
4632 state.buffer_fill(cid, buffered);
4633
4634 let accepted = make_status_report_with_quantity(
4635 Some("O-HER-4154-BUF"),
4636 new_voi,
4637 OrderStatus::Accepted,
4638 Quantity::from("0.835"),
4639 );
4640 let corrective = handle_execution_report(
4641 ExecutionReport::Order(accepted),
4642 &state,
4643 &emitter,
4644 &ws_client,
4645 &make_http_client(),
4646 &mut pending_cloids,
4647 UnixNanos::default(),
4648 );
4649
4650 let _ = drain_events(&mut rx);
4651 let (_, _, request) =
4652 corrective.expect("buffered fill drained before compute must still queue a corrective");
4653 assert_eq!(request.size, "0.835".parse::<Decimal>().unwrap());
4654 }
4655
4656 #[rstest]
4660 fn test_cancel_replace_no_corrective_when_filled_equals_target() {
4661 let ws_client = make_ws_client();
4662 let (emitter, mut rx) = test_emitter();
4663 let state = WsDispatchState::new();
4664 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4665
4666 let cid = ClientOrderId::from("O-HER-4154-EXACT");
4667 let target_total = Quantity::from("1.000");
4668
4669 let mut context = test_context(cid);
4670 context.quantity = target_total;
4671 state.register_context(context);
4672 state.insert_accepted(cid);
4673 state.record_venue_order_id(cid, VenueOrderId::new("445117664938"));
4674 state.mark_pending_modify(cid, VenueOrderId::new("445117664938"), target_total);
4675 state.stash_modify_request(cid, limit_request(Decimal::from(1)));
4676 state.record_filled_qty(cid, target_total);
4677
4678 let accepted = make_status_report_with_quantity(
4679 Some("O-HER-4154-EXACT"),
4680 "445117686214",
4681 OrderStatus::Accepted,
4682 target_total,
4683 );
4684 let corrective = handle_execution_report(
4685 ExecutionReport::Order(accepted),
4686 &state,
4687 &emitter,
4688 &ws_client,
4689 &make_http_client(),
4690 &mut pending_cloids,
4691 UnixNanos::default(),
4692 );
4693
4694 let _ = drain_events(&mut rx);
4695 assert!(corrective.is_none());
4696 assert!(state.pending_modify(&cid).is_none());
4697 }
4698
4699 #[rstest]
4702 fn test_cancel_replace_chains_second_corrective_reduce() {
4703 let ws_client = make_ws_client();
4704 let (emitter, mut rx) = test_emitter();
4705 let state = WsDispatchState::new();
4706 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4707
4708 let cid = ClientOrderId::from("O-HER-4154-CHAIN");
4709 let target_total = Quantity::from("1.000");
4710 let voi3 = "445117699999";
4711
4712 let mut context = test_context(cid);
4715 context.quantity = target_total;
4716 state.register_context(context);
4717 state.insert_accepted(cid);
4718 state.record_venue_order_id(cid, VenueOrderId::new("445117686214"));
4719 state.mark_pending_modify(cid, VenueOrderId::new("445117686214"), target_total);
4720 state.stash_modify_request(cid, limit_request("0.835".parse::<Decimal>().unwrap()));
4721 state.record_filled_qty(cid, Quantity::from("0.465"));
4723
4724 let accepted = make_status_report_with_quantity(
4725 Some("O-HER-4154-CHAIN"),
4726 voi3,
4727 OrderStatus::Accepted,
4728 Quantity::from("0.535"),
4729 );
4730 let corrective = handle_execution_report(
4731 ExecutionReport::Order(accepted),
4732 &state,
4733 &emitter,
4734 &ws_client,
4735 &make_http_client(),
4736 &mut pending_cloids,
4737 UnixNanos::default(),
4738 );
4739
4740 let _ = drain_events(&mut rx);
4741 let (_, oid, request) =
4742 corrective.expect("a further in-flight fill must chain another corrective");
4743 assert_eq!(oid, 445_117_699_999);
4744 assert_eq!(request.size, "0.535".parse::<Decimal>().unwrap());
4745 assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new(voi3)));
4746 }
4747
4748 #[rstest]
4749 fn test_handle_execution_report_external_terminal_evicts_cloid() {
4750 let ws_client = make_ws_client();
4754 let (emitter, mut rx) = test_emitter();
4755 let state = WsDispatchState::new();
4756 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4757
4758 let cid = ClientOrderId::from("O-HER-EXT");
4759 ws_client.cache_cloid_mapping(cloid_for("O-HER-EXT"), cid);
4760
4761 let report = make_status_report(Some("O-HER-EXT"), "v-ext", OrderStatus::Canceled);
4762 handle_execution_report(
4763 ExecutionReport::Order(report),
4764 &state,
4765 &emitter,
4766 &ws_client,
4767 &make_http_client(),
4768 &mut pending_cloids,
4769 UnixNanos::default(),
4770 );
4771
4772 let events = drain_events(&mut rx);
4773 assert_eq!(events.len(), 1);
4774 assert!(
4775 matches!(events[0], ExecutionEvent::Report(_)),
4776 "external terminal report should forward to the engine as a report",
4777 );
4778 assert_eq!(ws_client.get_cloid_mapping(&cloid_for("O-HER-EXT")), None);
4779 }
4780
4781 #[rstest]
4782 fn test_handle_execution_report_open_status_preserves_cloid() {
4783 let ws_client = make_ws_client();
4785 let (emitter, _rx) = test_emitter();
4786 let state = WsDispatchState::new();
4787 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4788
4789 let cid = ClientOrderId::from("O-HER-OPEN");
4790 state.register_context(test_context(cid));
4791 ws_client.cache_cloid_mapping(cloid_for("O-HER-OPEN"), cid);
4792
4793 let report = make_status_report(Some("O-HER-OPEN"), "v-open", OrderStatus::Accepted);
4794 handle_execution_report(
4795 ExecutionReport::Order(report),
4796 &state,
4797 &emitter,
4798 &ws_client,
4799 &make_http_client(),
4800 &mut pending_cloids,
4801 UnixNanos::default(),
4802 );
4803
4804 assert_eq!(
4806 ws_client.get_cloid_mapping(&cloid_for("O-HER-OPEN")),
4807 Some(cid)
4808 );
4809 }
4810
4811 #[rstest]
4812 fn test_handle_execution_report_tracked_accepted_emits_typed_event() {
4813 let ws_client = make_ws_client();
4817 let (emitter, mut rx) = test_emitter();
4818 let state = WsDispatchState::new();
4819 let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
4820
4821 let cid = ClientOrderId::from("O-HER-ACC");
4822 state.register_context(test_context(cid));
4823 ws_client.cache_cloid_mapping(cloid_for("O-HER-ACC"), cid);
4824
4825 let report = make_status_report(Some("O-HER-ACC"), "v-acc", OrderStatus::Accepted);
4826 handle_execution_report(
4827 ExecutionReport::Order(report),
4828 &state,
4829 &emitter,
4830 &ws_client,
4831 &make_http_client(),
4832 &mut pending_cloids,
4833 UnixNanos::default(),
4834 );
4835
4836 let events = drain_events(&mut rx);
4837 assert_eq!(events.len(), 1);
4838 assert!(
4839 matches!(events[0], ExecutionEvent::Order(OrderEventAny::Accepted(_))),
4840 "tracked accepted should route through the typed-event path",
4841 );
4842 assert_eq!(
4844 ws_client.get_cloid_mapping(&cloid_for("O-HER-ACC")),
4845 Some(cid)
4846 );
4847 }
4848
4849 fn outcome_limit_order(id: &str, reduce_only: bool) -> OrderAny {
4850 outcome_limit_order_full(id, reduce_only, false, TimeInForce::Gtc)
4851 }
4852
4853 fn outcome_limit_order_full(
4854 id: &str,
4855 reduce_only: bool,
4856 post_only: bool,
4857 time_in_force: TimeInForce,
4858 ) -> OrderAny {
4859 OrderAny::Limit(LimitOrder::new(
4860 TraderId::from("TESTER-001"),
4861 StrategyId::from("S-001"),
4862 InstrumentId::from("1-YES-OUTCOME.HYPERLIQUID"),
4863 ClientOrderId::from(id),
4864 OrderSide::Buy,
4865 Quantity::from("1"),
4866 Price::from("0.5000"),
4867 time_in_force,
4868 None,
4869 post_only,
4870 reduce_only,
4871 false,
4872 None,
4873 None,
4874 None,
4875 Some(ContingencyType::NoContingency),
4876 None,
4877 None,
4878 None,
4879 None,
4880 None,
4881 None,
4882 None,
4883 Default::default(),
4884 Default::default(),
4885 ))
4886 }
4887
4888 fn outcome_stop_order(id: &str) -> OrderAny {
4889 OrderAny::StopMarket(StopMarketOrder::new(
4890 TraderId::from("TESTER-001"),
4891 StrategyId::from("S-001"),
4892 InstrumentId::from("1-YES-OUTCOME.HYPERLIQUID"),
4893 ClientOrderId::from(id),
4894 OrderSide::Sell,
4895 Quantity::from("1"),
4896 Price::from("0.4000"),
4897 TriggerType::LastPrice,
4898 TimeInForce::Gtc,
4899 None,
4900 false,
4901 false,
4902 None,
4903 None,
4904 None,
4905 Some(ContingencyType::NoContingency),
4906 None,
4907 None,
4908 None,
4909 None,
4910 None,
4911 None,
4912 None,
4913 Default::default(),
4914 Default::default(),
4915 ))
4916 }
4917
4918 fn perp_with_unsupported_symbol(id: &str) -> OrderAny {
4919 OrderAny::Limit(LimitOrder::new(
4920 TraderId::from("TESTER-001"),
4921 StrategyId::from("S-001"),
4922 InstrumentId::from("BTC-USD-FOO.HYPERLIQUID"),
4923 ClientOrderId::from(id),
4924 OrderSide::Buy,
4925 Quantity::from("1"),
4926 Price::from("100.0"),
4927 TimeInForce::Gtc,
4928 None,
4929 false,
4930 false,
4931 false,
4932 None,
4933 None,
4934 None,
4935 Some(ContingencyType::NoContingency),
4936 None,
4937 None,
4938 None,
4939 None,
4940 None,
4941 None,
4942 None,
4943 Default::default(),
4944 Default::default(),
4945 ))
4946 }
4947
4948 #[rstest]
4949 fn test_validate_accepts_perp_limit_order() {
4950 let order = limit_order(
4951 "O-VAL-PERP",
4952 false,
4953 ContingencyType::NoContingency,
4954 None,
4955 None,
4956 );
4957 validate_order_for_hyperliquid(&order).unwrap();
4958 }
4959
4960 #[rstest]
4961 #[case::gtc_post_only(true, TimeInForce::Gtc)]
4962 #[case::gtc_taker(false, TimeInForce::Gtc)]
4963 #[case::ioc_post_only(true, TimeInForce::Ioc)]
4964 #[case::ioc_taker(false, TimeInForce::Ioc)]
4965 fn test_validate_accepts_outcome_limit_order(
4966 #[case] post_only: bool,
4967 #[case] time_in_force: TimeInForce,
4968 ) {
4969 let order = outcome_limit_order_full(
4970 "O-VAL-OUTCOME",
4971 false,
4972 post_only,
4973 time_in_force,
4974 );
4975 validate_order_for_hyperliquid(&order).unwrap();
4976 }
4977
4978 #[rstest]
4979 fn test_validate_rejects_outcome_reduce_only() {
4980 let order = outcome_limit_order("O-VAL-RO", true);
4981 let err = validate_order_for_hyperliquid(&order).unwrap_err();
4982 assert!(
4983 err.to_string().contains("Reduce-only is not supported"),
4984 "unexpected error: {err}",
4985 );
4986 }
4987
4988 #[rstest]
4989 fn test_validate_rejects_outcome_trigger_order() {
4990 let order = outcome_stop_order("O-VAL-TRIG");
4991 let err = validate_order_for_hyperliquid(&order).unwrap_err();
4992 assert!(
4993 err.to_string()
4994 .contains("Trigger order types are not supported"),
4995 "unexpected error: {err}",
4996 );
4997 }
4998
4999 #[rstest]
5000 fn test_validate_rejects_unsupported_symbol_suffix() {
5001 let order = perp_with_unsupported_symbol("O-VAL-BAD");
5002 let err = validate_order_for_hyperliquid(&order).unwrap_err();
5003 assert!(
5004 err.to_string()
5005 .contains("Unsupported instrument symbol format"),
5006 "unexpected error: {err}",
5007 );
5008 }
5009}