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