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