Skip to main content

nautilus_hyperliquid/
execution.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live execution client implementation for the Hyperliquid adapter.
17
18use 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    /// Returns a reference to the configuration.
122    pub fn config(&self) -> &HyperliquidExecutionClientConfig {
123        &self.config
124    }
125
126    /// Returns a reference to the shared WebSocket dispatch state.
127    ///
128    /// Exposes the context map, pending-modify markers, and cached venue
129    /// order ids used by the two-tier dispatch contract. The state is
130    /// read-write via an [`Arc`]; callers must not mutate it directly, but
131    /// it is useful for inspection in tests and for live debugging.
132    #[must_use]
133    pub fn ws_dispatch_state(&self) -> &Arc<WsDispatchState> {
134        &self.ws_dispatch_state
135    }
136
137    /// Returns `true` when every background task spawned via `spawn_task`
138    /// has completed.
139    ///
140    /// Used in tests to wait for submit / modify / cancel action round-trips
141    /// that fire on the runtime to finish before asserting on dispatch
142    /// state, avoiding bare `sleep` calls when a negative condition needs
143    /// to be checked after the spawned work is done.
144    #[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        // Market orders need a limit price derived from the cached quote,
194        // leaving the conversion's zero placeholder when none is cached.
195        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    // Hyperliquid accepts only base-denominated sizes on the wire, so a quote
214    // amount must convert (or deny) before it can reach the venue.
215    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    /// Creates a new [`HyperliquidExecutionClient`].
368    ///
369    /// # Errors
370    ///
371    /// Returns an error if either the HTTP or WebSocket client fail to construct.
372    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        // Emit even when both sides are empty so the account registers for
510        // await_account_registered on unfunded wallets.
511        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    /// Registers an instrument published by the data client so a market listed
762    /// after this client bootstrapped becomes submittable without a restart.
763    ///
764    /// The HTTP client takes the venue asset index from the instrument's `info`
765    /// map. WebSocket submissions sign through `&self.http_client`, so that one
766    /// write serves both submission paths.
767    fn on_instrument(&mut self, instrument: InstrumentAny) {
768        // this is the only step that makes a market listed after our bootstrap
769        // submittable, so it is traceable rather than silent
770        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        // Complete venue conversion before emitting OrderSubmitted
840        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                    // Deny MARKET orders without a cached quote, matching
1022                    // the single-order path.
1023                    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        // A bracket whose entry failed conversion must not submit its
1054        // children: they would rest at the venue as orphan exits.
1055        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        // A modify replaces a base-denominated size, and a quote target cannot
1202        // be reconciled against the venue's base-denominated fill reports.
1203        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        // Hyperliquid modify is cancel-replace; subtract filled to avoid overfill.
1273        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        // Build base request from cached order (derives slippage-adjusted
1308        // limit for trigger-market types like StopMarket/MarketIfTouched)
1309        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                        // Trigger changed but no explicit price: re-derive the
1328                        // slippage-adjusted limit from the new trigger
1329                        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        // Register on the command thread so concurrent tasks cannot stash a request
1398        // onto another modify intent. Failed admission clears only this generation.
1399        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            // The promotion uses this request to reduce a replacement after an in-flight fill
1406            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            // Search open orders by cloid first so modify/cancel-replace
1815            // resolves to the live replacement rather than a stale cached oid.
1816            // Request errors here are logged, not propagated, so a transient
1817            // frontendOpenOrders failure does not abort the whole query.
1818            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        // Post-WS setup: if any step fails, tear down WS before returning
1923        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        // Search open orders by cloid first when supplied. Hyperliquid modify
1995        // produces a new venue oid while preserving cloid, so a cached oid can
1996        // point at the canceled leg rather than the live replacement.
1997        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                            // A failed cloid probe is not a "not found": with no
2043                            // oid fallback available the lookup must fail closed.
2044                            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        // The venue bounds both history endpoints to their most recent entries;
2184        // a saturated response may be truncated, so coverage is not provable
2185        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        // Apply lookback filter to fills only (positions are current state,
2219        // and open orders must always be included for correct reconciliation)
2220        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        // Must match REST queries; mismatch silently drops fills on agent wallets
2272        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            // Cloids for external / untracked orders that reach a terminal
2314            // state: we evict their mapping immediately so long-running
2315            // sessions do not leak. Tracked orders clear their own cloid
2316            // mapping from the dispatch `cleanup_terminal` path below.
2317            //
2318            // For a tracked order that hits a status-only `FILLED` marker
2319            // without an accompanying fill, we defer the cloid cleanup until
2320            // the matching `FillReport` arrives so partial fills do not lose
2321            // their `client_order_id` link. The bounded FIFO cache keeps
2322            // orphaned entries from growing unbounded.
2323            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
2628// During a tracked cancel-replace the cached venue_order_id is still the old
2629// leg, so only its `Canceled` must be dropped (it would wrongly terminate the
2630// live order); a late `Filled` or any other status is forwarded for recovery.
2631fn 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
2949/// Registers an order's context in the dispatch state so its subsequent
2950/// WebSocket lifecycle can route through the typed-event path.
2951///
2952/// Quote-quantity orders submit a quote amount (e.g. 100 USD) but the venue
2953/// reports fills in base units. Comparing those two when deciding whether an
2954/// order is fully filled would leave the order stuck "open" forever, so they
2955/// flow through the untracked path and the engine reconciles them from
2956/// status reports instead.
2957fn 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
2994/// Validates that an order is acceptable for submission to Hyperliquid.
2995///
2996/// Checks symbol format, order type support, and HIP-4-specific restrictions
2997/// (no reduce-only, no trigger order types on outcome side tokens).
2998///
2999/// # Errors
3000///
3001/// Returns an error describing the first validation failure encountered.
3002pub 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    // HIP-4 outcomes are fully-collateralized side tokens with no margin,
3026    // funding, or trigger machinery. Reject features that don't apply.
3027    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
3529/// Routes a single execution report through the two-tier dispatch.
3530///
3531/// For tracked orders this emits typed `OrderEventAny` events via the
3532/// dispatch module; external / untracked orders fall back to the raw report
3533/// so the engine can reconcile. Cloid-mapping cleanup is handled here so
3534/// long-running sessions do not leak mapping entries.
3535fn 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            // Cloid cleanup:
3557            //
3558            // * `Skip` (stale cancel leg of a cancel-replace, cancel-before-accept
3559            //   race, or replay after terminal): leave the mapping intact. The
3560            //   still-open replacement order depends on it for subsequent events,
3561            //   and a genuinely terminal replay had its mapping evicted earlier.
3562            // * `Tracked` + status-only FILLED marker: defer the eviction until
3563            //   the matching `FillReport` lands so the partial fill preceding it
3564            //   keeps its client-order-id link.
3565            // * `Tracked` non-marker terminal and `External` terminal: evict now
3566            //   so long-running sessions do not leak cloid mappings.
3567            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            // Hand any queued corrective reduce to the loop to post; this
3582            // cache-free task cannot rebuild the order spec itself.
3583            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            // Skip cleanup while a cancel-replace fill is buffered; the
3599            // replacement ACCEPTED still needs to resolve the cloid (GH-3972).
3600            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
3617/// Posts a corrective reduce queued by the cancel-replace promotion.
3618///
3619/// Runs on the runtime (not the WS receive loop) so the post does not block
3620/// event processing. Keeps the re-armed pending-modify marker only while the
3621/// reduce may still be live (a clean ack, or a transport failure the WS may
3622/// reconcile); clears it otherwise so a stale marker cannot suppress a later
3623/// real `CANCELED(oid)` as a cancel-before-accept leg.
3624fn 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` does not connect, so this is a
3791        // cheap unit-test shim that still exercises the real `cloid_cache`
3792        // mapping APIs used by `handle_execution_report`.
3793        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    // Matches the order built by `limit_order_with_flags` so registration
3828    // tests can assert the stored context by equality.
3829    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        // A market listed after the client bootstrapped arrives through the
3852        // message bus, and must become submittable without a restart.
3853        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,  // expire_time
4000            false, // post_only
4001            reduce_only,
4002            false, // quote_quantity
4003            None,  // display_qty
4004            None,  // emulation_trigger
4005            None,  // trigger_instrument_id
4006            contingency,
4007            None, // order_list_id
4008            linked_ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()),
4009            parent_id.map(ClientOrderId::from),
4010            None, // exec_algorithm_id
4011            None, // exec_algorithm_params
4012            None, // exec_spawn_id
4013            None, // tags
4014            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, // expire_time
4037            reduce_only,
4038            false, // quote_quantity
4039            None,  // display_qty
4040            None,  // emulation_trigger
4041            None,  // trigger_instrument_id
4042            contingency,
4043            None, // order_list_id
4044            linked_ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()),
4045            parent_id.map(ClientOrderId::from),
4046            None, // exec_algorithm_id
4047            None, // exec_algorithm_params
4048            None, // exec_spawn_id
4049            None, // tags
4050            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        // Quote-quantity orders flow through the untracked path so the engine
4376        // reconciles them from status reports; registering would make the
4377        // cumulative-fill comparison mismatch base-unit fills against the
4378        // quote-unit tracked quantity and leave the order stuck "open".
4379        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        // Regression guard for GH-3827: when the dispatch returns Skip (e.g.
4389        // the stale cancel leg of a cancel-replace), the cloid mapping must
4390        // stay in place so the still-open replacement order can still be
4391        // resolved by subsequent events.
4392        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        // A tracked CANCELED that reaches a genuine terminal state should
4426        // emit OrderCanceled and evict the cloid mapping so long-running
4427        // sessions do not leak.
4428        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        // The status-only FILLED marker defers the cloid eviction to the
4584        // pending cache; the matching FillReport emits OrderFilled and then
4585        // evicts the cloid mapping as part of the deferred-cleanup path.
4586        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    /// GH-4270: when a status-only `FILLED` marker arrives before the
4636    /// replacement fill and the `ACCEPTED(new_voi)` is dropped, the fill itself
4637    /// promotes the binding (OrderUpdated then OrderFilled) and, being terminal
4638    /// and no longer buffered, completes the deferred cloid eviction.
4639    #[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    /// After a partial fill, the cancel-replace `OrderUpdated` must carry
4708    /// the user's absolute total, not the venue's remaining-only view.
4709    #[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        // Terminal-fill detection uses the context's absolute target quantity
4759        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    /// Without a target-qty marker (e.g. reconcile-driven modifies), the
4773    /// promotion falls back to `report.quantity`.
4774    #[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    /// A partial fill landing mid-modify leaves the cancel-replace replacement
4832    /// oversized (sized at the full target). The promotion must queue a
4833    /// corrective reduce to `target - filled` and re-arm the marker.
4834    #[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        // Marker re-armed on the new voi so the corrective's own cancel leg
4889        // is suppressed and a further in-flight fill chains another reduce.
4890        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    /// GH-4270: when the replacement ACCEPTED is dropped, a fill on the new leg
4895    /// promotes the binding and (parity with the ACCEPTED path) queues a corrective
4896    /// reduce when an earlier old-leg fill left the replacement oversized.
4897    #[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    /// Without an in-flight fill the replacement is correctly sized, so the
4954    /// promotion must not queue a corrective reduce and must clear the marker.
4955    #[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    /// A fill buffered during the in-flight cancel-replace is drained before
4997    /// the corrective is computed, so the corrective must size from the
4998    /// post-drain cumulative. A pre-drain read would see 0 filled and skip it.
4999    #[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    /// When an in-flight fill brings cumulative filled to exactly the target,
5049    /// the remaining is zero, so no corrective (a reduce-to-zero is not valid)
5050    /// must be queued and the marker is cleared.
5051    #[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    /// A further in-flight fill during the corrective's own modify chains
5092    /// another reduce: the second promotion sizes from the new cumulative.
5093    #[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        // State after the first corrective: marker re-armed on the prior
5105        // replacement, stashed request reduced to 0.835, 0.165 already filled.
5106        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        // External (untracked) terminal reports forward to the engine via
5142        // send_order_status_report and immediately evict the cloid mapping
5143        // so the client does not leak mappings for orders it does not own.
5144        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        // A tracked open ACCEPTED must flow through the typed-event path,
5203        // NOT the raw report fallback. Catches a mutation that swaps the
5204        // branch polarity inside `handle_execution_report`.
5205        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            /* reduce_only */ 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}