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