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