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