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