use std::{
sync::{Arc, Mutex},
time::{Duration, Instant},
};
use ahash::AHashMap;
use anyhow::Context;
use async_trait::async_trait;
use nautilus_common::{
cache::fifo::FifoCache,
clients::ExecutionClient,
live::{runner::get_exec_event_sender, runtime::get_runtime, task::TaskHandles},
messages::execution::{
BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
GenerateOrderStatusReport, GenerateOrderStatusReports, GeneratePositionStatusReports,
ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
},
};
use nautilus_core::{
MUTEX_POISONED, Params, UUID4, UnixNanos,
time::{AtomicTime, get_atomic_clock_realtime},
};
use nautilus_live::{ExecutionClientCore, ExecutionEventEmitter};
use nautilus_model::{
accounts::AccountAny,
enums::{AccountType, OmsType, OrderSide, OrderStatus, OrderType},
identifiers::{
AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, Venue, VenueOrderId,
},
orders::{Order, any::OrderAny},
reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
types::{AccountBalance, MarginBalance, Quantity},
};
#[derive(Debug, Clone)]
struct StagedBracketChild {
order: OrderAny,
request: HyperliquidExecPlaceOrderRequest,
}
#[derive(Debug, Default)]
struct StagedBracketState {
children_by_parent: AHashMap<ClientOrderId, Vec<StagedBracketChild>>,
active_children: AHashMap<ClientOrderId, StagedBracketChild>,
active_siblings: AHashMap<ClientOrderId, ClientOrderId>,
}
impl StagedBracketState {
fn stage(&mut self, parent_id: ClientOrderId, children: Vec<StagedBracketChild>) {
self.children_by_parent.insert(parent_id, children);
}
fn activate(&mut self, parent_id: &ClientOrderId) -> Option<Vec<StagedBracketChild>> {
let children = self.children_by_parent.remove(parent_id)?;
self.track_active(&children);
Some(children)
}
fn restore_active(&mut self, children: &[StagedBracketChild]) {
self.track_active(children);
}
fn track_active(&mut self, children: &[StagedBracketChild]) {
let child_ids = children
.iter()
.map(|child| child.order.client_order_id())
.collect::<Vec<_>>();
for child in children {
let child_id = child.order.client_order_id();
if let Some(sibling_id) = child
.order
.linked_order_ids()
.and_then(|ids| ids.iter().find(|id| child_ids.contains(id)))
{
self.active_siblings.insert(child_id, *sibling_id);
}
self.active_children.insert(child_id, child.clone());
}
}
fn contains_parent(&self, parent_id: &ClientOrderId) -> bool {
self.children_by_parent.contains_key(parent_id)
}
fn cancel_child(&mut self, child_id: &ClientOrderId) -> Option<OrderAny> {
let parent_id = self
.children_by_parent
.iter()
.find_map(|(parent_id, children)| {
children
.iter()
.any(|child| child.order.client_order_id() == *child_id)
.then_some(*parent_id)
})?;
let children = self.children_by_parent.get_mut(&parent_id)?;
let index = children
.iter()
.position(|child| child.order.client_order_id() == *child_id)?;
let child = children.remove(index);
if children.is_empty() {
self.children_by_parent.remove(&parent_id);
}
Some(child.order)
}
fn cancel_for_parent(&mut self, parent_id: &ClientOrderId) -> Vec<OrderAny> {
self.children_by_parent
.remove(parent_id)
.map(|children| children.into_iter().map(|child| child.order).collect())
.unwrap_or_default()
}
fn take_active_sibling(
&mut self,
client_order_id: &ClientOrderId,
) -> Option<StagedBracketChild> {
self.active_children.remove(client_order_id);
let sibling_id = self.active_siblings.remove(client_order_id)?;
self.active_siblings.remove(&sibling_id);
self.active_children.remove(&sibling_id)
}
fn active_sibling(&self, client_order_id: &ClientOrderId) -> Option<StagedBracketChild> {
self.active_siblings
.get(client_order_id)
.and_then(|sibling_id| self.active_children.get(sibling_id))
.cloned()
}
}
use tokio::task::JoinHandle;
use ustr::Ustr;
use crate::{
account::resolve_execution_account_address,
common::{
consts::{
HYPERLIQUID_BUILDER_APPROVAL_DOCS_URL, HYPERLIQUID_BUILDER_FEE_NOT_APPROVED,
HYPERLIQUID_POST_ONLY_WOULD_MATCH, HYPERLIQUID_VENUE,
},
credential::Secrets,
enums::HyperliquidProductType,
parse::{
clamp_price_to_precision, derive_limit_from_trigger, derive_market_order_price,
extract_error_message, extract_inner_error, extract_inner_errors, normalize_price,
order_to_hyperliquid_request_with_asset_and_cloid,
parse_combined_account_balances_and_margins, round_to_sig_figs,
},
},
config::HyperliquidExecClientConfig,
http::{
client::HyperliquidHttpClient,
models::{
ClearinghouseState, Cloid, HyperliquidExecAction, HyperliquidExecCancelByCloidRequest,
HyperliquidExecCancelOrderRequest, HyperliquidExecGrouping,
HyperliquidExecModifyOrderRequest, HyperliquidExecModifyTarget,
HyperliquidExecOrderKind, HyperliquidExecPlaceOrderRequest, HyperliquidExecTpSl,
SpotClearinghouseState,
},
parse::derive_outcome_settlements,
},
outcome_settlement::{OutcomeSettlementTracker, build_settlement_fills},
websocket::{
ExecutionReport, NautilusWsMessage,
client::HyperliquidWebSocketClient,
dispatch::{
DispatchOutcome, OrderIdentity, WsDispatchState, dispatch_order_event,
dispatch_order_fill, promote_replacement_from_query,
},
},
};
#[derive(Debug)]
pub struct HyperliquidExecutionClient {
core: ExecutionClientCore,
clock: &'static AtomicTime,
config: HyperliquidExecClientConfig,
emitter: ExecutionEventEmitter,
http_client: HyperliquidHttpClient,
ws_client: HyperliquidWebSocketClient,
pending_tasks: TaskHandles,
ws_stream_handle: Option<JoinHandle<()>>,
settlement_poll_handle: Option<JoinHandle<()>>,
ws_dispatch_state: Arc<WsDispatchState>,
staged_brackets: Arc<Mutex<StagedBracketState>>,
outcome_settlement_tracker: Arc<Mutex<OutcomeSettlementTracker>>,
}
impl HyperliquidExecutionClient {
pub fn config(&self) -> &HyperliquidExecClientConfig {
&self.config
}
#[must_use]
pub fn ws_dispatch_state(&self) -> &Arc<WsDispatchState> {
&self.ws_dispatch_state
}
#[allow(
clippy::missing_panics_doc,
reason = "pending_tasks mutex poisoning is not expected"
)]
#[must_use]
pub fn pending_tasks_all_finished(&self) -> bool {
self.pending_tasks.all_finished()
}
fn resolve_slippage_bps(&self, params: Option<&Params>) -> u32 {
params
.and_then(|p| p.get_u64("market_order_slippage_bps"))
.map_or(self.config.market_order_slippage_bps, |v| v as u32)
}
fn validate_order_submission(&self, order: &OrderAny) -> anyhow::Result<()> {
validate_order_for_hyperliquid(order)
}
fn order_request(
&self,
order: &OrderAny,
slippage_bps: u32,
) -> anyhow::Result<HyperliquidExecPlaceOrderRequest> {
validate_order_for_hyperliquid(order)?;
let symbol = order.instrument_id().symbol.inner();
let asset = self
.http_client
.get_asset_index_for_symbol(symbol)
.with_context(|| format!("Asset index not found for {symbol}"))?;
let price_decimals = self
.http_client
.get_price_precision_for_symbol(symbol)
.unwrap_or(2);
let cloid = self
.http_client
.get_or_generate_client_order_id_cloid(order.client_order_id());
let mut request = order_to_hyperliquid_request_with_asset_and_cloid(
order,
asset,
price_decimals,
self.config.normalize_prices,
slippage_bps,
None,
)?;
request.cloid = Some(cloid);
Ok(request)
}
fn restore_staged_brackets(&self) -> Vec<ClientOrderId> {
let order_lists = self
.core
.cache()
.order_lists(Some(&self.core.venue), None, None, None)
.into_iter()
.cloned()
.collect::<Vec<_>>();
let mut ready_parent_ids = Vec::new();
for order_list in order_lists {
let orders = {
let cache = self.core.cache();
order_list
.client_order_ids
.iter()
.filter_map(|client_order_id| {
cache.order(client_order_id).map(|order| order.clone())
})
.collect::<Vec<_>>()
};
if orders.len() != order_list.client_order_ids.len()
|| determine_order_list_grouping(&orders) != HyperliquidExecGrouping::NormalTpsl
{
continue;
}
let (mut orders, mut requests) = match orders
.iter()
.map(|order| self.order_request(order, self.config.market_order_slippage_bps))
.collect::<anyhow::Result<Vec<_>>>()
{
Ok(requests) => order_normal_tpsl_submission(
orders,
requests,
HyperliquidExecGrouping::NormalTpsl,
),
Err(e) => {
log::warn!("Cannot restore staged bracket {}: {e}", order_list.id,);
continue;
}
};
let parent = orders.remove(0);
let parent_request = requests.remove(0);
let parent_id = parent.client_order_id();
let (staged_children, active_children): (Vec<_>, Vec<_>) = orders
.drain(..)
.zip(requests.drain(..))
.filter(|(order, _)| order.is_active_local())
.map(|(order, request)| StagedBracketChild { order, request })
.partition(|child| child.order.status() == OrderStatus::Initialized);
if (staged_children.is_empty() && active_children.is_empty())
|| (!parent.is_open() && parent.filled_qty().raw == 0)
|| self
.staged_brackets
.lock()
.expect(MUTEX_POISONED)
.contains_parent(&parent_id)
{
continue;
}
self.restore_order_identity(&parent, &parent_request);
for child in &active_children {
self.restore_order_identity(&child.order, &child.request);
}
let has_staged_children = !staged_children.is_empty();
let mut state = self.staged_brackets.lock().expect(MUTEX_POISONED);
if has_staged_children {
state.stage(parent_id, staged_children);
}
state.restore_active(&active_children);
drop(state);
if has_staged_children && parent.filled_qty().raw > 0 {
ready_parent_ids.push(parent_id);
}
}
if !ready_parent_ids.is_empty() {
log::info!(
"Restored {} staged bracket parent(s) with prior fills",
ready_parent_ids.len(),
);
}
ready_parent_ids
}
fn restore_order_identity(&self, order: &OrderAny, request: &HyperliquidExecPlaceOrderRequest) {
let client_order_id = order.client_order_id();
let cloid = request.cloid.expect("order conversion must set a CLOID");
self.http_client
.cache_client_order_id_cloid(client_order_id, cloid);
self.ws_client
.cache_cloid_mapping(Ustr::from(&cloid.to_hex()), client_order_id);
self.ws_dispatch_state.register_identity(
client_order_id,
OrderIdentity {
strategy_id: order.strategy_id(),
instrument_id: order.instrument_id(),
order_side: order.order_side(),
order_type: order.order_type(),
quantity: order.quantity(),
price: order.price(),
},
);
if let Some(venue_order_id) = order.venue_order_id() {
self.ws_dispatch_state
.record_venue_order_id(client_order_id, venue_order_id);
self.ws_dispatch_state.insert_accepted(client_order_id);
}
}
pub fn new(
core: ExecutionClientCore,
config: HyperliquidExecClientConfig,
) -> anyhow::Result<Self> {
let secrets = Secrets::resolve(
config.private_key.as_deref(),
config.vault_address.as_deref(),
config.environment,
)
.context("Hyperliquid execution client requires private key")?;
let account_address = resolve_execution_account_address(
config.private_key.as_deref(),
config.vault_address.as_deref(),
config.account_address.as_deref(),
config.environment,
)?;
let mut http_client = HyperliquidHttpClient::with_secrets(
&secrets,
config.http_timeout_secs,
config.proxy_url.clone(),
)
.context("failed to create Hyperliquid HTTP client")?;
http_client.set_account_id(core.account_id);
http_client.set_account_address(account_address);
http_client.set_normalize_prices(config.normalize_prices);
http_client.set_market_order_slippage_bps(config.market_order_slippage_bps);
http_client.set_include_builder_attribution(config.include_builder_attribution);
if let Some(url) = &config.base_url_http {
http_client.set_base_info_url(url.clone());
}
if let Some(url) = &config.base_url_exchange {
http_client.set_base_exchange_url(url.clone());
}
let ws_url = config.base_url_ws.clone();
let mut ws_client = HyperliquidWebSocketClient::new(
ws_url,
config.environment,
Some(core.account_id),
config.transport_backend,
config.proxy_url.clone(),
);
ws_client.set_post_timeout(Duration::from_secs(config.ws_post_timeout_secs));
let clock = get_atomic_clock_realtime();
let emitter = ExecutionEventEmitter::new(
clock,
core.trader_id,
core.account_id,
AccountType::Margin,
None,
);
Ok(Self {
core,
clock,
config,
emitter,
http_client,
ws_client,
pending_tasks: TaskHandles::default(),
ws_stream_handle: None,
settlement_poll_handle: None,
ws_dispatch_state: Arc::new(WsDispatchState::new()),
staged_brackets: Arc::new(Mutex::new(StagedBracketState::default())),
outcome_settlement_tracker: Arc::new(Mutex::new(OutcomeSettlementTracker::new())),
})
}
fn register_order_identity(&self, order: &OrderAny) {
register_order_identity_into(&self.ws_dispatch_state, order);
}
async fn ensure_instruments_initialized_async(&self) -> anyhow::Result<()> {
if self.core.instruments_initialized() {
return Ok(());
}
let instruments = self
.http_client
.request_instruments()
.await
.context("failed to request Hyperliquid instruments")?;
if instruments.is_empty() {
log::warn!(
"Instrument bootstrap yielded no instruments; WebSocket submissions may fail"
);
} else {
log::debug!("Initialized {} instruments", instruments.len());
for instrument in &instruments {
self.http_client.cache_instrument(instrument);
}
}
self.core.set_instruments_initialized();
Ok(())
}
async fn refresh_account_state(&self) -> anyhow::Result<()> {
let account_address = self.get_account_address()?;
let (perp_state, spot_state) = self
.fetch_combined_clearinghouse_state(&account_address)
.await?;
log::debug!(
"Received clearinghouse state: cross_margin_summary={:?}, asset_positions={}, spot_balances={}",
perp_state.cross_margin_summary,
perp_state.asset_positions.len(),
spot_state.balances.len(),
);
let (balances, margins) =
parse_combined_account_balances_and_margins(&perp_state, &spot_state)
.context("failed to parse combined account balances and margins")?;
let ts_event = self.clock.get_time_ns();
self.emitter
.emit_account_state(balances, margins, true, ts_event);
log::debug!("Account state updated successfully");
Ok(())
}
async fn fetch_combined_clearinghouse_state(
&self,
account_address: &str,
) -> anyhow::Result<(ClearinghouseState, SpotClearinghouseState)> {
let perp_json = self
.http_client
.info_clearinghouse_state(account_address)
.await
.context("failed to fetch clearinghouse state")?;
let perp_state: ClearinghouseState = serde_json::from_value(perp_json)
.context("failed to deserialize clearinghouse state")?;
let spot_json = self
.http_client
.info_spot_clearinghouse_state(account_address)
.await
.context("failed to fetch spot clearinghouse state")?;
let spot_state: SpotClearinghouseState = serde_json::from_value(spot_json)
.context("failed to deserialize spot clearinghouse state")?;
Ok((perp_state, spot_state))
}
async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
let account_id = self.core.account_id;
if self.core.cache().account(&account_id).is_some() {
log::info!("Account {account_id} registered");
return Ok(());
}
let start = Instant::now();
let timeout = Duration::from_secs_f64(timeout_secs);
let interval = Duration::from_millis(10);
loop {
tokio::time::sleep(interval).await;
if self.core.cache().account(&account_id).is_some() {
log::info!("Account {account_id} registered");
return Ok(());
}
if start.elapsed() >= timeout {
anyhow::bail!(
"Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
);
}
}
}
fn get_account_address(&self) -> anyhow::Result<String> {
self.http_client
.get_account_address()
.context("failed to get account address from HTTP client")
}
fn spawn_task<F>(&self, description: &'static str, fut: F)
where
F: std::future::Future<Output = anyhow::Result<()>> + Send + 'static,
{
let runtime = get_runtime();
let handle = runtime.spawn(async move {
if let Err(e) = fut.await {
log::warn!("{description} failed: {e:?}");
}
});
self.pending_tasks.push(handle);
}
fn start_outcome_settlement_poll(&mut self) -> anyhow::Result<()> {
let poll_secs = self.config.outcome_settlement_poll_secs;
if poll_secs == 0 {
log::debug!("Outcome settlement polling disabled by config");
return Ok(());
}
let http_client = self.http_client.clone();
let emitter = self.emitter.clone();
let tracker = self.outcome_settlement_tracker.clone();
let account_id = self.core.account_id;
let account_address = self.get_account_address()?;
let clock = self.clock;
let handle = get_runtime().spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(poll_secs));
interval.tick().await;
loop {
interval.tick().await;
let meta = match http_client.get_outcome_meta().await {
Ok(meta) => meta,
Err(e) => {
log::warn!("Outcome meta poll failed: {e}");
continue;
}
};
let settlements = derive_outcome_settlements(&meta);
if settlements.is_empty() {
continue;
}
let spot_json = match http_client
.info_spot_clearinghouse_state(&account_address)
.await
{
Ok(value) => value,
Err(e) => {
log::warn!("Settlement dispatch skipped: spot state fetch failed: {e}");
continue;
}
};
let spot_state: SpotClearinghouseState = match serde_json::from_value(spot_json) {
Ok(state) => state,
Err(e) => {
log::warn!("Settlement dispatch skipped: spot state parse failed: {e}");
continue;
}
};
let ts = clock.get_time_ns();
let fills = {
let mut guard = tracker.lock().expect(MUTEX_POISONED);
build_settlement_fills(&settlements, &spot_state, &mut guard, account_id, ts)
};
for fill in fills {
log::debug!(
"Dispatching outcome settlement fill: instrument={}, price={}, qty={}",
fill.instrument_id,
fill.last_px,
fill.last_qty,
);
emitter.send_fill_report(fill);
}
}
});
if let Some(previous) = self.settlement_poll_handle.replace(handle) {
previous.abort();
}
Ok(())
}
fn abort_pending_tasks(&self) {
self.pending_tasks.abort_all();
}
}
#[async_trait(?Send)]
impl ExecutionClient for HyperliquidExecutionClient {
fn is_connected(&self) -> bool {
self.core.is_connected()
}
fn client_id(&self) -> ClientId {
self.core.client_id
}
fn account_id(&self) -> AccountId {
self.core.account_id
}
fn venue(&self) -> Venue {
*HYPERLIQUID_VENUE
}
fn oms_type(&self) -> OmsType {
self.core.oms_type
}
fn get_account(&self) -> Option<AccountAny> {
self.core.cache().account_owned(&self.core.account_id)
}
fn generate_account_state(
&self,
balances: Vec<AccountBalance>,
margins: Vec<MarginBalance>,
reported: bool,
ts_event: UnixNanos,
) -> anyhow::Result<()> {
self.emitter
.emit_account_state(balances, margins, reported, ts_event);
Ok(())
}
fn start(&mut self) -> anyhow::Result<()> {
if self.core.is_started() {
return Ok(());
}
let sender = get_exec_event_sender();
self.emitter.set_sender(sender);
self.core.set_started();
log::info!(
"Started: client_id={}, account_id={}, environment={:?}, vault_address={:?}, proxy_url={:?}",
self.core.client_id,
self.core.account_id,
self.config.environment,
self.config.vault_address,
self.config.proxy_url,
);
Ok(())
}
fn stop(&mut self) -> anyhow::Result<()> {
if self.core.is_stopped() {
return Ok(());
}
log::info!("Stopping Hyperliquid execution client");
if let Some(handle) = self.ws_stream_handle.take() {
handle.abort();
}
if let Some(handle) = self.settlement_poll_handle.take() {
handle.abort();
}
self.abort_pending_tasks();
self.ws_client.abort();
self.core.set_disconnected();
self.core.set_stopped();
log::info!("Hyperliquid execution client stopped");
Ok(())
}
fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
if order.is_closed() {
log::warn!("Cannot submit closed order {}", order.client_order_id());
return Ok(());
}
if let Err(e) = self.validate_order_submission(&order) {
self.emitter
.emit_order_denied(&order, &format!("Validation failed: {e}"));
return Err(e);
}
let http_client = self.http_client.clone();
let symbol = order.instrument_id().symbol.inner();
let asset = match http_client.get_asset_index_for_symbol(symbol) {
Some(a) => a,
None => {
self.emitter
.emit_order_denied(&order, &format!("Asset index not found for {symbol}"));
return Ok(());
}
};
let price_decimals = http_client
.get_price_precision_for_symbol(symbol)
.unwrap_or(2);
let slippage_bps = self.resolve_slippage_bps(cmd.params.as_ref());
let mut hyperliquid_order = match order_to_hyperliquid_request_with_asset_and_cloid(
&order,
asset,
price_decimals,
self.config.normalize_prices,
slippage_bps,
None,
) {
Ok(req) => req,
Err(e) => {
self.emitter
.emit_order_denied(&order, &format!("Order conversion failed: {e}"));
return Ok(());
}
};
let cloid = http_client.get_or_generate_client_order_id_cloid(order.client_order_id());
hyperliquid_order.cloid = Some(cloid);
if order.order_type() == OrderType::Market {
let instrument_id = order.instrument_id();
let cache = self.core.cache();
match cache.quote(&instrument_id) {
Some(quote) => {
let is_buy = order.order_side() == OrderSide::Buy;
hyperliquid_order.price =
derive_market_order_price(quote, is_buy, price_decimals, slippage_bps);
}
None => {
self.emitter.emit_order_denied(
&order,
&format!(
"No cached quote for {instrument_id}: \
subscribe to quote data before submitting market orders"
),
);
return Ok(());
}
}
}
log::debug!(
"Submitting order: id={}, type={:?}, side={:?}, price={}, size={}, kind={:?}",
order.client_order_id(),
order.order_type(),
order.order_side(),
hyperliquid_order.price,
hyperliquid_order.size,
hyperliquid_order.kind,
);
let cloid = hyperliquid_order
.cloid
.expect("order conversion must set a CLOID");
self.http_client
.cache_client_order_id_cloid(order.client_order_id(), cloid);
self.ws_client
.cache_cloid_mapping(Ustr::from(&cloid.to_hex()), order.client_order_id());
self.register_order_identity(&order);
self.emitter.emit_order_submitted(&order);
let emitter = self.emitter.clone();
let clock = self.clock;
let ws_client = self.ws_client.clone();
let cloid_hex = Ustr::from(&cloid.to_hex());
let dispatch_state = self.ws_dispatch_state.clone();
let builder = self.http_client.builder_attribution();
self.spawn_task("submit_order", async move {
let action = HyperliquidExecAction::Order {
orders: vec![hyperliquid_order],
grouping: HyperliquidExecGrouping::Na,
builder,
};
let rejection_route =
PostRejectionRoute::new(&emitter, &ws_client, &http_client, dispatch_state.clone());
match ws_client.post_action_exec(&http_client, &action).await {
Ok(response) => {
if response.is_ok() {
if let Some(inner_error) = extract_inner_error(&response) {
log::warn!("Order submission rejected by exchange: {inner_error}");
let ts = clock.get_time_ns();
rejection_route.emit_once(&order, &inner_error, ts, &cloid_hex);
} else {
log::debug!("Order submitted successfully: {response:?}");
}
} else {
let error_msg = extract_error_message(&response);
log::warn!("Order submission rejected by exchange: {error_msg}");
let ts = clock.get_time_ns();
rejection_route.emit_once(&order, &error_msg, ts, &cloid_hex);
}
}
Err(e) => {
log::error!("Order submission WebSocket post request failed: {e}");
}
}
rejection_route.resolve_without_post_rejection(&order, clock.get_time_ns(), &cloid_hex);
Ok(())
});
Ok(())
}
fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
log::debug!(
"Submitting order list with {} orders",
cmd.order_list.client_order_ids.len()
);
let http_client = self.http_client.clone();
let slippage_bps = self.resolve_slippage_bps(cmd.params.as_ref());
let orders = self.core.get_orders_for_list(&cmd.order_list)?;
let mut valid_orders = Vec::new();
let mut hyperliquid_orders = Vec::new();
for order in &orders {
match self.order_request(order, slippage_bps) {
Ok(request) => {
hyperliquid_orders.push(request);
valid_orders.push(order.clone());
}
Err(e) => {
self.emitter
.emit_order_denied(order, &format!("Order conversion failed: {e}"));
}
}
}
if valid_orders.is_empty() {
log::warn!("No valid orders to submit in order list");
return Ok(());
}
let grouping = determine_order_list_grouping(&valid_orders);
log::debug!("Order list grouping: {grouping:?}");
let (mut valid_orders, mut hyperliquid_orders) =
order_normal_tpsl_submission(valid_orders, hyperliquid_orders, grouping);
let submission_grouping = if grouping == HyperliquidExecGrouping::NormalTpsl {
let parent = valid_orders.remove(0);
let parent_request = hyperliquid_orders.remove(0);
let children = valid_orders
.drain(..)
.zip(hyperliquid_orders.drain(..))
.map(|(order, request)| StagedBracketChild { order, request })
.collect();
self.staged_brackets
.lock()
.expect(MUTEX_POISONED)
.stage(parent.client_order_id(), children);
valid_orders.push(parent);
hyperliquid_orders.push(parent_request);
HyperliquidExecGrouping::Na
} else {
grouping
};
for (order, request) in valid_orders.iter().zip(hyperliquid_orders.iter()) {
let cloid = request.cloid.expect("order conversion must set a CLOID");
self.http_client
.cache_client_order_id_cloid(order.client_order_id(), cloid);
self.ws_client
.cache_cloid_mapping(Ustr::from(&cloid.to_hex()), order.client_order_id());
self.register_order_identity(order);
self.emitter.emit_order_submitted(order);
}
let emitter = self.emitter.clone();
let clock = self.clock;
let ws_client = self.ws_client.clone();
let dispatch_state = self.ws_dispatch_state.clone();
let staged_brackets = self.staged_brackets.clone();
let builder = self.http_client.builder_attribution();
self.spawn_task("submit_order_list", async move {
post_order_batch(
"Order list",
valid_orders,
hyperliquid_orders,
submission_grouping,
builder,
&emitter,
&ws_client,
&http_client,
dispatch_state,
staged_brackets,
clock,
)
.await;
Ok(())
});
Ok(())
}
fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
log::debug!("Modifying order: {cmd:?}");
let client_order_id = cmd.client_order_id;
let venue_order_id = cmd
.venue_order_id
.or_else(|| self.core.cache().venue_order_id(&client_order_id).copied());
let order = match self.core.cache().order(&client_order_id).map(|o| o.clone()) {
Some(o) => o,
None => {
let reason = "order not found in cache";
log::warn!("Cannot modify order {client_order_id}: {reason}");
self.emitter.emit_order_modify_rejected_event(
cmd.strategy_id,
cmd.instrument_id,
client_order_id,
venue_order_id,
reason,
self.clock.get_time_ns(),
);
return Ok(());
}
};
let http_client = self.http_client.clone();
let symbol = cmd.instrument_id.symbol.inner();
let should_normalize = self.config.normalize_prices;
let slippage_bps = self.resolve_slippage_bps(cmd.params.as_ref());
let modify_target = match http_client.unique_cached_client_order_id_cloid(&client_order_id)
{
Some(cloid) => HyperliquidExecModifyTarget::Cloid(cloid),
None => {
let Some(venue_order_id) = venue_order_id.as_ref() else {
let reason = "venue_order_id or unique cached CLOID is required for modify";
log::warn!("Cannot modify order {client_order_id}: {reason}");
self.emitter.emit_order_modify_rejected_event(
cmd.strategy_id,
cmd.instrument_id,
client_order_id,
None,
reason,
self.clock.get_time_ns(),
);
return Ok(());
};
match HyperliquidExecModifyTarget::from_venue_order_id(venue_order_id) {
Ok(target) => target,
Err(e) => {
let reason =
format!("Failed to parse venue_order_id '{venue_order_id}': {e}");
log::warn!("{reason}");
self.emitter.emit_order_modify_rejected_event(
cmd.strategy_id,
cmd.instrument_id,
client_order_id,
Some(*venue_order_id),
&reason,
self.clock.get_time_ns(),
);
return Ok(());
}
}
}
};
let old_venue_order_id = venue_order_id.filter(|id| id.as_str().parse::<u64>().is_ok());
if matches!(modify_target, HyperliquidExecModifyTarget::Cloid(_))
&& old_venue_order_id.is_none()
{
let reason = "cached venue_order_id is required for CLOID modify";
log::warn!("Cannot modify order {client_order_id}: {reason}");
self.emitter.emit_order_modify_rejected_event(
cmd.strategy_id,
cmd.instrument_id,
client_order_id,
venue_order_id,
reason,
self.clock.get_time_ns(),
);
return Ok(());
}
let target_total_qty = cmd.quantity.unwrap_or(order.quantity());
let filled_qty = order.filled_qty();
if target_total_qty <= filled_qty {
let reason =
format!("modify quantity {target_total_qty} not greater than filled {filled_qty}",);
log::warn!("Cannot modify order {}: {reason}", cmd.client_order_id);
self.emitter.emit_order_modify_rejected_event(
cmd.strategy_id,
cmd.instrument_id,
client_order_id,
venue_order_id,
&reason,
self.clock.get_time_ns(),
);
return Ok(());
}
let quantity = target_total_qty - filled_qty;
let price_decimals = http_client
.get_price_precision_for_symbol(symbol)
.unwrap_or(2);
let asset = match http_client.get_asset_index_for_symbol(symbol) {
Some(a) => a,
None => {
log::warn!(
"Asset index not found for symbol {symbol}, ensure instruments are loaded",
);
return Ok(());
}
};
let mut hyperliquid_order = match order_to_hyperliquid_request_with_asset_and_cloid(
&order,
asset,
price_decimals,
should_normalize,
slippage_bps,
None,
) {
Ok(mut req) => {
if let Some(p) = cmd.price.or(order.price()) {
let price_dec = p.as_decimal();
req.price = if should_normalize {
normalize_price(price_dec, price_decimals).normalize()
} else {
price_dec.normalize()
};
} else if let Some(tp) = cmd.trigger_price {
let is_buy = order.order_side() == OrderSide::Buy;
let base = tp.as_decimal().normalize();
let derived = derive_limit_from_trigger(base, is_buy, slippage_bps);
let sig_rounded = round_to_sig_figs(derived, 5);
req.price =
clamp_price_to_precision(sig_rounded, price_decimals, is_buy).normalize();
}
req.size = quantity.as_decimal().normalize();
if let (Some(tp), HyperliquidExecOrderKind::Trigger { trigger }) =
(cmd.trigger_price, &mut req.kind)
{
let tp_dec = tp.as_decimal();
trigger.trigger_px = if should_normalize {
normalize_price(tp_dec, price_decimals).normalize()
} else {
tp_dec.normalize()
};
}
req
}
Err(e) => {
log::warn!("Order conversion failed for modify: {e}");
return Ok(());
}
};
let cached_cloid_before_modify = http_client.cached_client_order_id_cloid(&client_order_id);
let cloid = http_client.get_or_generate_client_order_id_cloid(order.client_order_id());
let generated_modify_cloid = cached_cloid_before_modify
.is_none()
.then_some((client_order_id, cloid));
hyperliquid_order.cloid = Some(cloid);
let dispatch_state = self.ws_dispatch_state.clone();
let ws_client = self.ws_client.clone();
if let Some(cloid) = hyperliquid_order.cloid {
http_client.cache_client_order_id_cloid(client_order_id, cloid);
ws_client.cache_cloid_mapping(Ustr::from(&cloid.to_hex()), client_order_id);
}
let modify_generation = old_venue_order_id.map(|old_venue_order_id| {
let generation = dispatch_state.mark_pending_modify(
client_order_id,
old_venue_order_id,
target_total_qty,
);
dispatch_state.stash_modify_request(client_order_id, hyperliquid_order.clone());
generation
});
self.spawn_task("modify_order", async move {
let action = HyperliquidExecAction::Modify {
modify: HyperliquidExecModifyOrderRequest {
oid: modify_target,
order: hyperliquid_order,
},
};
match ws_client.post_action_exec(&http_client, &action).await {
Ok(response) => {
if response.is_ok() {
if let Some(inner_error) = extract_inner_error(&response) {
log::warn!("Order modification rejected by exchange: {inner_error}");
if let Some(generation) = modify_generation {
dispatch_state
.clear_modify_generation(&client_order_id, generation);
}
remove_generated_modify_cloid(
&http_client,
&ws_client,
generated_modify_cloid,
);
} else {
log::debug!("Order modified successfully: {response:?}");
}
} else {
let error_msg = extract_error_message(&response);
log::warn!("Order modification rejected by exchange: {error_msg}");
if let Some(generation) = modify_generation {
dispatch_state.clear_modify_generation(&client_order_id, generation);
}
remove_generated_modify_cloid(
&http_client,
&ws_client,
generated_modify_cloid,
);
}
}
Err(e) => {
if e.is_transport_error() {
log::warn!(
"Order modification transport failure for {client_order_id}: {e}; \
awaiting WS reconciliation",
);
} else {
log::warn!("Order modification WebSocket post request failed: {e}");
if let Some(generation) = modify_generation {
dispatch_state.clear_modify_generation(&client_order_id, generation);
}
remove_generated_modify_cloid(
&http_client,
&ws_client,
generated_modify_cloid,
);
}
}
}
Ok(())
});
Ok(())
}
fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
log::debug!("Cancelling order: {cmd:?}");
if let Some(order) = self
.staged_brackets
.lock()
.expect(MUTEX_POISONED)
.cancel_child(&cmd.client_order_id)
{
self.emitter
.emit_order_canceled(&order, None, self.clock.get_time_ns());
return Ok(());
}
let http_client = self.http_client.clone();
let emitter = self.emitter.clone();
let clock = self.clock;
let client_order_id = cmd.client_order_id;
let strategy_id = cmd.strategy_id;
let instrument_id = cmd.instrument_id;
let venue_order_id = cmd.venue_order_id;
let symbol = cmd.instrument_id.symbol.inner();
let ws_client = self.ws_client.clone();
let fast = can_fast_cancel_order(
self.core
.cache()
.order(&client_order_id)
.as_ref()
.map(|order| order.order_type()),
)
.then_some(true);
self.spawn_task("cancel_order", async move {
let asset = match http_client.get_asset_index_for_symbol(symbol) {
Some(a) => a,
None => {
log::warn!(
"Local cancel validation failed for {client_order_id}: Asset index not found for symbol {symbol}"
);
return Ok(());
}
};
let action =
if let Some(cloid) = http_client.cached_client_order_id_cloid(&client_order_id) {
HyperliquidExecAction::CancelByCloid {
cancels: vec![HyperliquidExecCancelByCloidRequest { asset, cloid }],
fast,
}
} else if let Some(venue_order_id) = venue_order_id {
match venue_order_id.as_str().parse::<u64>() {
Ok(oid) => HyperliquidExecAction::Cancel {
cancels: vec![HyperliquidExecCancelOrderRequest { asset, oid }],
fast,
},
Err(_) => {
log::warn!(
"Local cancel validation failed for {client_order_id}: Invalid venue order ID format"
);
return Ok(());
}
}
} else {
let cloid = http_client.get_or_generate_client_order_id_cloid(client_order_id);
HyperliquidExecAction::CancelByCloid {
cancels: vec![HyperliquidExecCancelByCloidRequest { asset, cloid }],
fast,
}
};
match ws_client.post_action_exec(&http_client, &action).await {
Ok(response) => {
if response.is_ok() {
if let Some(inner_error) = extract_inner_error(&response) {
emitter.emit_order_cancel_rejected_event(
strategy_id,
instrument_id,
client_order_id,
venue_order_id,
&inner_error,
clock.get_time_ns(),
);
} else {
log::debug!("Order cancelled successfully: {response:?}");
}
} else {
let error_msg = extract_error_message(&response);
log::warn!(
"Cancel failed without per-order result for {client_order_id}, awaiting WS reconciliation: {error_msg}"
);
}
}
Err(e) => {
if e.is_transport_error() {
log::warn!(
"Cancel transport failure for {client_order_id}: {e}; \
awaiting WS reconciliation",
);
} else {
log::warn!(
"Ambiguous cancel failure for {client_order_id}, awaiting WS reconciliation: {e}"
);
}
}
}
Ok(())
});
Ok(())
}
fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
log::debug!("Cancelling all orders: {cmd:?}");
let cache = self.core.cache();
let open_orders = cache.orders_open(
Some(&self.core.venue),
Some(&cmd.instrument_id),
None,
None,
Some(cmd.order_side),
);
if open_orders.is_empty() {
log::debug!("No open orders to cancel for {:?}", cmd.instrument_id);
return Ok(());
}
let symbol = cmd.instrument_id.symbol.inner();
let instrument_id = cmd.instrument_id;
let strategy_id = cmd.strategy_id;
let entries: Vec<CancelEntry> = open_orders
.iter()
.map(|o| CancelEntry {
strategy_id,
instrument_id,
client_order_id: o.client_order_id(),
venue_order_id: o.venue_order_id(),
symbol,
fast: can_fast_cancel_order(Some(o.order_type())),
})
.collect();
let http_client = self.http_client.clone();
let emitter = self.emitter.clone();
let clock = self.clock;
let ws_client = self.ws_client.clone();
self.spawn_task("cancel_all_orders", async move {
let asset = match http_client.get_asset_index_for_symbol(symbol) {
Some(a) => a,
None => {
log::warn!(
"Local cancel-all validation failed: Asset index not found for symbol {symbol}"
);
return Ok(());
}
};
let mut cancel_dispatch = CancelDispatch::new();
for entry in &entries {
cancel_dispatch.push(entry, asset, &http_client);
}
if cancel_dispatch.is_empty() {
return Ok(());
}
submit_cancel_dispatch(
"Cancel-all",
cancel_dispatch,
&ws_client,
&http_client,
&emitter,
clock,
)
.await;
Ok(())
});
Ok(())
}
fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
log::debug!("Batch cancelling orders: {cmd:?}");
if cmd.cancels.is_empty() {
log::debug!("No orders to cancel in batch");
return Ok(());
}
let cache = self.core.cache();
let entries: Vec<CancelEntry> = cmd
.cancels
.iter()
.map(|c| CancelEntry {
strategy_id: c.strategy_id,
instrument_id: c.instrument_id,
client_order_id: c.client_order_id,
venue_order_id: c.venue_order_id,
symbol: c.instrument_id.symbol.inner(),
fast: can_fast_cancel_order(
cache
.order(&c.client_order_id)
.as_ref()
.map(|order| order.order_type()),
),
})
.collect();
let http_client = self.http_client.clone();
let emitter = self.emitter.clone();
let clock = self.clock;
let ws_client = self.ws_client.clone();
self.spawn_task("batch_cancel_orders", async move {
let mut cancel_dispatch = CancelDispatch::new();
for entry in &entries {
let asset = match http_client.get_asset_index_for_symbol(entry.symbol) {
Some(a) => a,
None => {
log::warn!(
"Local batch cancel validation failed for {}: Asset index not found for symbol {}",
entry.client_order_id,
entry.symbol,
);
continue;
}
};
cancel_dispatch.push(entry, asset, &http_client);
}
if cancel_dispatch.is_empty() {
log::warn!("No valid cancel requests in batch");
return Ok(());
}
submit_cancel_dispatch(
"Batch cancel",
cancel_dispatch,
&ws_client,
&http_client,
&emitter,
clock,
)
.await;
Ok(())
});
Ok(())
}
fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
let http_client = self.http_client.clone();
let account_address = self.get_account_address()?;
let emitter = self.emitter.clone();
let clock = self.clock;
self.spawn_task("query_account", async move {
let perp_json = http_client
.info_clearinghouse_state(&account_address)
.await
.context("failed to fetch clearinghouse state")?;
let perp_state: ClearinghouseState = serde_json::from_value(perp_json)
.context("failed to deserialize clearinghouse state")?;
let spot_json = http_client
.info_spot_clearinghouse_state(&account_address)
.await
.context("failed to fetch spot clearinghouse state")?;
let spot_state: SpotClearinghouseState = serde_json::from_value(spot_json)
.context("failed to deserialize spot clearinghouse state")?;
let (balances, margins) =
parse_combined_account_balances_and_margins(&perp_state, &spot_state)
.context("failed to parse combined account balances and margins")?;
let ts_event = clock.get_time_ns();
emitter.emit_account_state(balances, margins, true, ts_event);
Ok(())
});
Ok(())
}
fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
log::debug!("Querying order: {cmd:?}");
let client_order_id = cmd.client_order_id;
let venue_order_id = match cmd.venue_order_id {
Some(voi) => Some(voi),
None => self.core.cache().venue_order_id(&client_order_id).copied(),
};
let account_address = self.get_account_address()?;
let http_client = self.http_client.clone();
let emitter = self.emitter.clone();
let dispatch_state = self.ws_dispatch_state.clone();
let clock = self.clock;
self.spawn_task("query_order", async move {
match http_client
.request_order_status_report_by_client_order_id(&account_address, &client_order_id)
.await
{
Ok(Some(report)) => {
promote_replacement_from_query(
&report,
&dispatch_state,
&emitter,
clock.get_time_ns(),
);
log::debug!("Queried order status for {client_order_id}");
emitter.send_order_status_report(report);
return Ok(());
}
Ok(None) => {}
Err(e) => {
log::warn!(
"Failed to query order status for {client_order_id}: {e}; falling back to oid lookup"
);
}
}
let Some(venue_order_id) = venue_order_id else {
log::debug!("No order status report found for {client_order_id}");
return Ok(());
};
let oid: u64 = match venue_order_id.as_str().parse() {
Ok(oid) => oid,
Err(e) => {
log::warn!("Failed to parse venue order ID {venue_order_id}: {e}");
return Ok(());
}
};
match http_client
.request_order_status_report(&account_address, oid)
.await
{
Ok(Some(report)) => {
if is_inflight_modify_old_leg_cancel(
&dispatch_state,
&client_order_id,
&report,
) {
log::debug!(
"Suppressing stale old-leg Canceled for {client_order_id}: modify in flight"
);
} else {
log::debug!("Queried order status for oid {oid}");
emitter.send_order_status_report(report);
}
}
Ok(None) => {
log::debug!("No order status report found for oid {oid}");
}
Err(e) => {
log::warn!("Failed to query order status for oid {oid}: {e}");
}
}
Ok(())
});
Ok(())
}
async fn connect(&mut self) -> anyhow::Result<()> {
if self.core.is_connected() {
return Ok(());
}
log::info!("Connecting Hyperliquid execution client");
self.ensure_instruments_initialized_async().await?;
let ready_bracket_parents = self.restore_staged_brackets();
self.start_ws_stream().await?;
let post_ws = async {
self.refresh_account_state().await?;
self.await_account_registered(30.0).await?;
Ok::<(), anyhow::Error>(())
};
if let Err(e) = post_ws.await {
log::warn!("Connect failed after WS started, tearing down: {e}");
let _ = self.ws_client.disconnect().await;
self.abort_pending_tasks();
return Err(e);
}
for parent_id in ready_bracket_parents {
if let Some(children) = self
.staged_brackets
.lock()
.expect(MUTEX_POISONED)
.activate(&parent_id)
{
spawn_staged_children(
children,
&self.emitter,
&self.ws_client,
&self.http_client,
self.ws_dispatch_state.clone(),
self.staged_brackets.clone(),
self.http_client.builder_attribution(),
self.clock,
);
}
}
if let Err(e) = self.start_outcome_settlement_poll() {
log::warn!("Outcome settlement polling not started: {e}");
}
self.core.set_connected();
log::info!("Connected: client_id={}", self.core.client_id);
Ok(())
}
async fn disconnect(&mut self) -> anyhow::Result<()> {
if self.core.is_disconnected() {
return Ok(());
}
log::info!("Disconnecting Hyperliquid execution client");
self.ws_client.disconnect().await?;
if let Some(handle) = self.settlement_poll_handle.take() {
handle.abort();
}
self.abort_pending_tasks();
self.core.set_disconnected();
log::info!("Disconnected: client_id={}", self.core.client_id);
Ok(())
}
async fn generate_order_status_report(
&self,
cmd: &GenerateOrderStatusReport,
) -> anyhow::Result<Option<OrderStatusReport>> {
let account_address = self.get_account_address()?;
if cmd.venue_order_id.is_none() && cmd.client_order_id.is_none() {
log::warn!(
"Cannot generate order status report without venue_order_id or client_order_id"
);
return Ok(None);
}
if let Some(client_order_id) = &cmd.client_order_id {
match self
.http_client
.request_order_status_report_by_client_order_id(&account_address, client_order_id)
.await
{
Ok(Some(report)) => {
promote_replacement_from_query(
&report,
&self.ws_dispatch_state,
&self.emitter,
self.clock.get_time_ns(),
);
log::debug!("Generated order status report for {client_order_id}");
return Ok(Some(report));
}
Ok(None) => {}
Err(e) => {
log::warn!(
"Failed to generate order status report for {client_order_id}: {e}; \
falling back to oid lookup"
);
}
}
}
let oid = match &cmd.venue_order_id {
Some(venue_order_id) => venue_order_id
.as_str()
.parse::<u64>()
.context("failed to parse venue_order_id as oid")?,
None => match &cmd.client_order_id {
Some(client_order_id) => {
let cached_oid: Option<u64> = self
.core
.cache()
.venue_order_id(client_order_id)
.and_then(|v| v.as_str().parse::<u64>().ok());
match cached_oid {
Some(oid) => oid,
None => {
log::debug!("No order status report found for {client_order_id}");
return Ok(None);
}
}
}
None => unreachable!("cmd must carry at least one identifier"),
},
};
let report = self
.http_client
.request_order_status_report(&account_address, oid)
.await
.context("failed to generate order status report")?;
if let Some(report) = &report
&& let Some(client_order_id) = &cmd.client_order_id
&& is_inflight_modify_old_leg_cancel(&self.ws_dispatch_state, client_order_id, report)
{
log::debug!(
"Suppressing stale old-leg Canceled for {client_order_id}: modify in flight"
);
return Ok(None);
}
if report.is_some() {
log::debug!("Generated order status report for oid {oid}");
} else {
log::debug!("No order status report found for oid {oid}");
}
Ok(report)
}
async fn generate_order_status_reports(
&self,
cmd: &GenerateOrderStatusReports,
) -> anyhow::Result<Vec<OrderStatusReport>> {
let account_address = self.get_account_address()?;
let reports = self
.http_client
.request_order_status_reports(&account_address, cmd.instrument_id)
.await
.context("failed to generate order status reports")?;
let reports = filter_order_status_reports_for_command(reports, cmd);
log::debug!("Generated {} order status reports", reports.len());
Ok(reports)
}
async fn generate_fill_reports(
&self,
cmd: GenerateFillReports,
) -> anyhow::Result<Vec<FillReport>> {
let account_address = self.get_account_address()?;
let reports = self
.http_client
.request_fill_reports(&account_address, cmd.instrument_id)
.await
.context("failed to generate fill reports")?;
let reports = if let (Some(start), Some(end)) = (cmd.start, cmd.end) {
reports
.into_iter()
.filter(|r| r.ts_event >= start && r.ts_event <= end)
.collect()
} else if let Some(start) = cmd.start {
reports
.into_iter()
.filter(|r| r.ts_event >= start)
.collect()
} else if let Some(end) = cmd.end {
reports.into_iter().filter(|r| r.ts_event <= end).collect()
} else {
reports
};
log::debug!("Generated {} fill reports", reports.len());
Ok(reports)
}
async fn generate_position_status_reports(
&self,
cmd: &GeneratePositionStatusReports,
) -> anyhow::Result<Vec<PositionStatusReport>> {
let account_address = self.get_account_address()?;
let reports = self
.http_client
.request_position_status_reports(&account_address, cmd.instrument_id)
.await
.context("failed to generate position status reports")?;
log::debug!("Generated {} position status reports", reports.len());
Ok(reports)
}
async fn generate_mass_status(
&self,
lookback_mins: Option<u64>,
) -> anyhow::Result<Option<ExecutionMassStatus>> {
let ts_init = self.clock.get_time_ns();
let order_cmd = GenerateOrderStatusReports::new(
UUID4::new(),
ts_init,
true, None,
None,
None,
None,
None,
);
let fill_cmd =
GenerateFillReports::new(UUID4::new(), ts_init, None, None, None, None, None, None);
let position_cmd =
GeneratePositionStatusReports::new(UUID4::new(), ts_init, None, None, None, None, None);
let mut order_reports = self.generate_order_status_reports(&order_cmd).await?;
let mut fill_reports = self.generate_fill_reports(fill_cmd).await?;
let position_reports = self.generate_position_status_reports(&position_cmd).await?;
if let Some(mins) = lookback_mins {
let cutoff_ns = ts_init
.as_u64()
.saturating_sub(mins.saturating_mul(60).saturating_mul(1_000_000_000));
let cutoff = UnixNanos::from(cutoff_ns);
fill_reports.retain(|r| r.ts_event >= cutoff);
}
if !fill_reports.is_empty() {
let account_address = self.get_account_address()?;
let filled_order_ids: ahash::AHashSet<_> = fill_reports
.iter()
.map(|report| report.venue_order_id)
.collect();
let open_order_ids: ahash::AHashSet<_> = order_reports
.iter()
.map(|report| report.venue_order_id)
.collect();
let mut historical_reports = self
.http_client
.request_historical_order_status_reports(&account_address, None)
.await
.context("failed to generate historical order status reports")?;
historical_reports.retain(|report| {
filled_order_ids.contains(&report.venue_order_id)
&& !open_order_ids.contains(&report.venue_order_id)
});
order_reports.extend(historical_reports);
}
let mut mass_status = ExecutionMassStatus::new(
self.core.client_id,
self.core.account_id,
self.core.venue,
ts_init,
None,
);
mass_status.add_order_reports(order_reports);
mass_status.add_fill_reports(fill_reports);
mass_status.add_position_reports(position_reports);
log::info!(
"Generated mass status: {} orders, {} fills, {} positions",
mass_status.order_reports().len(),
mass_status.fill_reports().len(),
mass_status.position_reports().len(),
);
Ok(Some(mass_status))
}
}
impl HyperliquidExecutionClient {
async fn start_ws_stream(&mut self) -> anyhow::Result<()> {
if self.ws_stream_handle.is_some() {
return Ok(());
}
let subscription_address = self.get_account_address()?;
let mut ws_client = self.ws_client.clone();
let instruments = self
.http_client
.request_instruments()
.await
.unwrap_or_default();
for instrument in instruments {
ws_client.cache_instrument(instrument);
}
ws_client.connect().await?;
ws_client
.subscribe_order_updates(&subscription_address)
.await?;
ws_client
.subscribe_user_events(&subscription_address)
.await?;
log::debug!("Subscribed to Hyperliquid execution updates for {subscription_address}");
if let Some(handle) = ws_client.take_task_handle() {
self.ws_client.set_task_handle(handle);
}
let emitter = self.emitter.clone();
let dispatch_state = self.ws_dispatch_state.clone();
let staged_brackets = self.staged_brackets.clone();
let http_client = self.http_client.clone();
let builder = self.http_client.builder_attribution();
let clock = self.clock;
let runtime = get_runtime();
let handle = runtime.spawn(async move {
let mut pending_filled_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
loop {
let event = ws_client.next_event().await;
match event {
Some(msg) => match msg {
NautilusWsMessage::ExecutionReports(reports) => {
for report in reports {
let staged_parent_fill = match &report {
ExecutionReport::Fill(report) => report.client_order_id,
ExecutionReport::Order(_) => None,
};
let staged_parent_terminal = match &report {
ExecutionReport::Order(report)
if matches!(
report.order_status,
OrderStatus::Canceled
| OrderStatus::Rejected
| OrderStatus::Expired
) =>
{
report.client_order_id.map(|client_order_id| {
(client_order_id, report.ts_last)
})
}
_ => None,
};
let active_child_terminal = match &report {
ExecutionReport::Order(report)
if matches!(
report.order_status,
OrderStatus::Filled
| OrderStatus::Canceled
| OrderStatus::Rejected
| OrderStatus::Expired
) =>
{
report.client_order_id
}
ExecutionReport::Fill(report) => {
report.client_order_id.filter(|client_order_id| {
let Some(identity) =
dispatch_state.lookup_identity(client_order_id)
else {
return false;
};
let previous = dispatch_state
.previous_filled_qty(client_order_id)
.unwrap_or_else(|| {
Quantity::zero(report.last_qty.precision)
});
previous + report.last_qty >= identity.quantity
})
}
_ => None,
};
let active_child_fill = match &report {
ExecutionReport::Fill(report) => {
report.client_order_id.and_then(|client_order_id| {
dispatch_state.lookup_identity(&client_order_id).map(
|identity| {
(
client_order_id,
dispatch_state
.previous_filled_qty(&client_order_id)
.unwrap_or_else(|| {
Quantity::zero(
report.last_qty.precision,
)
}),
identity.quantity,
)
},
)
})
}
ExecutionReport::Order(_) => None,
};
if let Some((cid, oid, order)) = handle_execution_report(
report,
&dispatch_state,
&emitter,
&ws_client,
&http_client,
&mut pending_filled_cloids,
clock.get_time_ns(),
) {
spawn_corrective_reduce(
&ws_client,
&http_client,
&dispatch_state,
cid,
oid,
order,
);
}
if let Some(parent_id) = staged_parent_fill
&& let Some(children) = staged_brackets
.lock()
.expect(MUTEX_POISONED)
.activate(&parent_id)
{
spawn_staged_children(
children,
&emitter,
&ws_client,
&http_client,
dispatch_state.clone(),
staged_brackets.clone(),
builder.clone(),
clock,
);
}
if let Some((parent_id, ts_event)) = staged_parent_terminal {
let children = staged_brackets
.lock()
.expect(MUTEX_POISONED)
.cancel_for_parent(&parent_id);
for child in children {
emitter.emit_order_canceled(&child, None, ts_event);
}
}
if let Some((client_order_id, previous, quantity)) =
active_child_fill
&& let Some(cumulative) =
dispatch_state.previous_filled_qty(&client_order_id)
&& cumulative > previous
&& cumulative < quantity
{
let sibling = staged_brackets
.lock()
.expect(MUTEX_POISONED)
.active_sibling(&client_order_id);
if let Some(sibling) = sibling {
spawn_active_sibling_resize(
sibling,
quantity - cumulative,
&emitter,
&ws_client,
&http_client,
&dispatch_state,
);
}
}
if let Some(client_order_id) = active_child_terminal {
let sibling = staged_brackets
.lock()
.expect(MUTEX_POISONED)
.take_active_sibling(&client_order_id);
if let Some(sibling) = sibling {
spawn_active_sibling_cancel(
sibling,
&emitter,
&ws_client,
&http_client,
&dispatch_state,
);
}
}
}
}
NautilusWsMessage::Reconnected => {}
NautilusWsMessage::Error(e) => {
log::warn!("WebSocket error: {e}");
}
NautilusWsMessage::Trades(_)
| NautilusWsMessage::Quote(_)
| NautilusWsMessage::Deltas(_)
| NautilusWsMessage::Depth10(_)
| NautilusWsMessage::Candle(_)
| NautilusWsMessage::MarkPrice(_)
| NautilusWsMessage::IndexPrice(_)
| NautilusWsMessage::FundingRate(_)
| NautilusWsMessage::CustomData(_) => {}
},
None => {
log::debug!("WebSocket next_event returned None, stream closed");
break;
}
}
}
});
self.ws_stream_handle = Some(handle);
log::debug!("Hyperliquid WebSocket execution stream started");
Ok(())
}
}
fn filter_order_status_reports_for_command(
reports: Vec<OrderStatusReport>,
cmd: &GenerateOrderStatusReports,
) -> Vec<OrderStatusReport> {
let reports = if cmd.open_only {
reports
.into_iter()
.filter(|r| r.order_status.is_open())
.collect()
} else {
reports
};
match (cmd.start, cmd.end) {
(Some(start), Some(end)) => reports
.into_iter()
.filter(|r| r.ts_last >= start && r.ts_last <= end)
.collect(),
(Some(start), None) => reports.into_iter().filter(|r| r.ts_last >= start).collect(),
(None, Some(end)) => reports.into_iter().filter(|r| r.ts_last <= end).collect(),
(None, None) => reports,
}
}
fn is_inflight_modify_old_leg_cancel(
dispatch_state: &WsDispatchState,
client_order_id: &ClientOrderId,
report: &OrderStatusReport,
) -> bool {
report.order_status == OrderStatus::Canceled
&& dispatch_state.pending_modify_contains_old(client_order_id, report.venue_order_id)
}
fn remove_generated_modify_cloid(
http_client: &HyperliquidHttpClient,
ws_client: &HyperliquidWebSocketClient,
generated_modify_cloid: Option<(ClientOrderId, Cloid)>,
) {
let Some((client_order_id, cloid)) = generated_modify_cloid else {
return;
};
if http_client.cached_client_order_id_cloid(&client_order_id) != Some(cloid) {
return;
}
let cloid_hex = Ustr::from(&cloid.to_hex());
ws_client.remove_cloid_mapping(&cloid_hex);
http_client.remove_client_order_id_cloid(&client_order_id);
}
#[derive(Clone)]
struct CancelEntry {
strategy_id: StrategyId,
instrument_id: InstrumentId,
client_order_id: ClientOrderId,
venue_order_id: Option<VenueOrderId>,
symbol: Ustr,
fast: bool,
}
struct CancelDispatch {
cloid_requests: Vec<(HyperliquidExecCancelByCloidRequest, CancelEntry)>,
oid_requests: Vec<(HyperliquidExecCancelOrderRequest, CancelEntry)>,
}
impl CancelDispatch {
fn new() -> Self {
Self {
cloid_requests: Vec::new(),
oid_requests: Vec::new(),
}
}
fn is_empty(&self) -> bool {
self.cloid_requests.is_empty() && self.oid_requests.is_empty()
}
fn push(&mut self, entry: &CancelEntry, asset: u32, http_client: &HyperliquidHttpClient) {
if let Some(cloid) = http_client.cached_client_order_id_cloid(&entry.client_order_id) {
self.cloid_requests.push((
HyperliquidExecCancelByCloidRequest { asset, cloid },
entry.clone(),
));
} else if let Some(venue_order_id) = entry.venue_order_id {
match venue_order_id.as_str().parse::<u64>() {
Ok(oid) => {
self.oid_requests.push((
HyperliquidExecCancelOrderRequest { asset, oid },
entry.clone(),
));
}
Err(_) => {
log::warn!(
"Local cancel validation failed for {}: Invalid venue order ID format",
entry.client_order_id,
);
}
}
} else {
let cloid = http_client.get_or_generate_client_order_id_cloid(entry.client_order_id);
self.cloid_requests.push((
HyperliquidExecCancelByCloidRequest { asset, cloid },
entry.clone(),
));
}
}
}
async fn submit_cancel_dispatch(
label: &str,
dispatch: CancelDispatch,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
emitter: &ExecutionEventEmitter,
clock: &'static AtomicTime,
) {
let CancelDispatch {
cloid_requests,
oid_requests,
} = dispatch;
let (fast_cloid_requests, fast_cloid_entries, cloid_requests, cloid_entries) =
split_fast_cancel_requests(cloid_requests);
if !fast_cloid_requests.is_empty() {
let action = HyperliquidExecAction::CancelByCloid {
cancels: fast_cloid_requests,
fast: Some(true),
};
submit_cancel_action(
label,
action,
&fast_cloid_entries,
ws_client,
http_client,
emitter,
clock,
)
.await;
}
if !cloid_requests.is_empty() {
let action = HyperliquidExecAction::CancelByCloid {
cancels: cloid_requests,
fast: None,
};
submit_cancel_action(
label,
action,
&cloid_entries,
ws_client,
http_client,
emitter,
clock,
)
.await;
}
let (fast_oid_requests, fast_oid_entries, oid_requests, oid_entries) =
split_fast_cancel_requests(oid_requests);
if !fast_oid_requests.is_empty() {
let action = HyperliquidExecAction::Cancel {
cancels: fast_oid_requests,
fast: Some(true),
};
submit_cancel_action(
label,
action,
&fast_oid_entries,
ws_client,
http_client,
emitter,
clock,
)
.await;
}
if !oid_requests.is_empty() {
let action = HyperliquidExecAction::Cancel {
cancels: oid_requests,
fast: None,
};
submit_cancel_action(
label,
action,
&oid_entries,
ws_client,
http_client,
emitter,
clock,
)
.await;
}
}
fn split_fast_cancel_requests<T>(
requests: Vec<(T, CancelEntry)>,
) -> (Vec<T>, Vec<CancelEntry>, Vec<T>, Vec<CancelEntry>) {
let mut fast_requests = Vec::new();
let mut fast_entries = Vec::new();
let mut requests_without_fast = Vec::new();
let mut entries_without_fast = Vec::new();
for (request, entry) in requests {
if entry.fast {
fast_requests.push(request);
fast_entries.push(entry);
} else {
requests_without_fast.push(request);
entries_without_fast.push(entry);
}
}
(
fast_requests,
fast_entries,
requests_without_fast,
entries_without_fast,
)
}
async fn submit_cancel_action(
label: &str,
action: HyperliquidExecAction,
sent_entries: &[CancelEntry],
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
emitter: &ExecutionEventEmitter,
clock: &'static AtomicTime,
) {
match ws_client.post_action_exec(http_client, &action).await {
Ok(response) => {
if response.is_ok() {
let inner_errors = extract_inner_errors(&response);
let ts = clock.get_time_ns();
if inner_errors.is_empty() {
log::debug!("{label} submitted successfully: {response:?}");
} else if let Some(reason) = cancel_status_count_mismatch_reason(
label,
sent_entries.len(),
inner_errors.len(),
) {
log::warn!("{reason}");
} else {
for (i, entry) in sent_entries.iter().enumerate() {
if let Some(Some(error_msg)) = inner_errors.get(i) {
log::warn!(
"Cancel for {} rejected by exchange: {error_msg}",
entry.client_order_id,
);
emitter.emit_order_cancel_rejected_event(
entry.strategy_id,
entry.instrument_id,
entry.client_order_id,
entry.venue_order_id,
error_msg,
ts,
);
}
}
}
} else {
let error_msg = extract_error_message(&response);
log::warn!(
"{label} failed without per-order results, awaiting WS reconciliation: {error_msg}"
);
}
}
Err(e) => {
if e.is_transport_error() {
log::warn!("{label} transport failure: {e}; awaiting WS reconciliation");
} else {
log::warn!("{label} ambiguous failure, awaiting WS reconciliation: {e}");
}
}
}
}
fn register_order_identity_into(state: &WsDispatchState, order: &OrderAny) {
if order.is_quote_quantity() {
return;
}
state.register_identity(
order.client_order_id(),
OrderIdentity {
strategy_id: order.strategy_id(),
instrument_id: order.instrument_id(),
order_side: order.order_side(),
order_type: order.order_type(),
quantity: order.quantity(),
price: order.price(),
},
);
state.mark_submission_pending(order.client_order_id());
}
fn order_normal_tpsl_submission(
orders: Vec<OrderAny>,
requests: Vec<HyperliquidExecPlaceOrderRequest>,
grouping: HyperliquidExecGrouping,
) -> (Vec<OrderAny>, Vec<HyperliquidExecPlaceOrderRequest>) {
if grouping != HyperliquidExecGrouping::NormalTpsl {
return (orders, requests);
}
let mut pairs: Vec<_> = orders.into_iter().zip(requests).collect();
pairs.sort_by_key(|(order, request)| {
if !order.is_reduce_only() {
0
} else if matches!(
&request.kind,
HyperliquidExecOrderKind::Trigger { trigger }
if trigger.tpsl == HyperliquidExecTpSl::Sl
) {
2
} else {
1
}
});
pairs.into_iter().unzip()
}
pub fn validate_order_for_hyperliquid(order: &OrderAny) -> anyhow::Result<()> {
let instrument_id = order.instrument_id();
let symbol = instrument_id.symbol.as_str();
let product_type = HyperliquidProductType::from_symbol(symbol).map_err(|_| {
anyhow::anyhow!(
"Unsupported instrument symbol format for Hyperliquid: {symbol} \
(expected -PERP, -SPOT, or HIP-4 outcome `{{N}}-{{YES|NO}}-OUTCOME`)"
)
})?;
match order.order_type() {
OrderType::Market
| OrderType::Limit
| OrderType::StopMarket
| OrderType::StopLimit
| OrderType::MarketIfTouched
| OrderType::LimitIfTouched => {}
_ => anyhow::bail!(
"Unsupported order type for Hyperliquid: {:?}",
order.order_type()
),
}
if product_type == HyperliquidProductType::Outcome {
if order.is_reduce_only() {
anyhow::bail!("Reduce-only is not supported for Hyperliquid HIP-4 outcomes: {symbol}");
}
if !matches!(order.order_type(), OrderType::Market | OrderType::Limit) {
anyhow::bail!(
"Trigger order types are not supported for Hyperliquid HIP-4 outcomes: \
{symbol} (received {:?})",
order.order_type()
);
}
}
if matches!(
order.order_type(),
OrderType::StopMarket
| OrderType::StopLimit
| OrderType::MarketIfTouched
| OrderType::LimitIfTouched
) && order.trigger_price().is_none()
{
anyhow::bail!(
"Conditional orders require a trigger price for Hyperliquid: {:?}",
order.order_type()
);
}
if matches!(
order.order_type(),
OrderType::Limit | OrderType::StopLimit | OrderType::LimitIfTouched
) && order.price().is_none()
{
anyhow::bail!(
"Limit orders require a limit price for Hyperliquid: {:?}",
order.order_type()
);
}
Ok(())
}
fn can_fast_cancel_order(order_type: Option<OrderType>) -> bool {
matches!(order_type, Some(OrderType::Market | OrderType::Limit))
}
fn cancel_status_count_mismatch_reason(
label: &str,
expected_count: usize,
actual_count: usize,
) -> Option<String> {
(actual_count != 0 && actual_count != expected_count).then(|| {
format!(
"{label} response status count mismatch: expected {expected_count}, received {actual_count}"
)
})
}
#[expect(clippy::too_many_arguments)]
async fn post_order_batch(
label: &str,
orders: Vec<OrderAny>,
requests: Vec<HyperliquidExecPlaceOrderRequest>,
grouping: HyperliquidExecGrouping,
builder: Option<crate::http::models::HyperliquidExecBuilderFee>,
emitter: &ExecutionEventEmitter,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
dispatch_state: Arc<WsDispatchState>,
staged_brackets: Arc<Mutex<StagedBracketState>>,
clock: &'static AtomicTime,
) {
let cloid_hexes: Vec<Ustr> = requests
.iter()
.map(|request| {
Ustr::from(
&request
.cloid
.expect("order conversion must set a CLOID")
.to_hex(),
)
})
.collect();
let action = HyperliquidExecAction::Order {
orders: requests,
grouping,
builder,
};
let rejection_route = PostRejectionRoute::with_staged_brackets(
emitter,
ws_client,
http_client,
dispatch_state,
staged_brackets,
);
match ws_client.post_action_exec(http_client, &action).await {
Ok(response) if response.is_ok() => {
let inner_errors = extract_inner_errors(&response);
let ts = clock.get_time_ns();
if inner_errors.len() == orders.len() {
for ((order, cloid_hex), error) in orders
.iter()
.zip(cloid_hexes.iter())
.zip(inner_errors.iter())
{
if let Some(error_msg) = error {
log::warn!(
"Order {} rejected by exchange: {error_msg}",
order.client_order_id(),
);
rejection_route.emit_once(order, error_msg, ts, cloid_hex);
}
}
} else if orders.len() > 1
&& inner_errors.len() == 1
&& let Some(error_msg) = inner_errors[0].as_ref()
{
log::warn!("{label} rejected by deterministic whole-batch validation: {error_msg}",);
for (order, cloid_hex) in orders.iter().zip(cloid_hexes.iter()) {
rejection_route.emit_once(order, error_msg, ts, cloid_hex);
}
} else if !inner_errors.is_empty() {
log::warn!(
"{label} returned {} statuses for {} orders; preserving unresolved identities \
for WebSocket or startup reconciliation",
inner_errors.len(),
orders.len(),
);
} else {
log::debug!("{label} submitted successfully: {response:?}");
}
}
Ok(response) => {
let error_msg = extract_error_message(&response);
log::warn!("{label} submission rejected by exchange: {error_msg}");
let ts = clock.get_time_ns();
for (order, cloid_hex) in orders.iter().zip(cloid_hexes.iter()) {
rejection_route.emit_once(order, &error_msg, ts, cloid_hex);
}
}
Err(e) => {
log::error!("{label} WebSocket post request failed: {e}");
}
}
let ts = clock.get_time_ns();
for (order, cloid_hex) in orders.iter().zip(cloid_hexes.iter()) {
rejection_route.resolve_without_post_rejection(order, ts, cloid_hex);
}
}
#[expect(clippy::too_many_arguments)]
fn spawn_staged_children(
children: Vec<StagedBracketChild>,
emitter: &ExecutionEventEmitter,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
dispatch_state: Arc<WsDispatchState>,
staged_brackets: Arc<Mutex<StagedBracketState>>,
builder: Option<crate::http::models::HyperliquidExecBuilderFee>,
clock: &'static AtomicTime,
) {
let (orders, requests): (Vec<_>, Vec<_>) = children
.into_iter()
.map(|child| (child.order, child.request))
.unzip();
for (order, request) in orders.iter().zip(requests.iter()) {
let cloid = request.cloid.expect("order conversion must set a CLOID");
http_client.cache_client_order_id_cloid(order.client_order_id(), cloid);
ws_client.cache_cloid_mapping(Ustr::from(&cloid.to_hex()), order.client_order_id());
register_order_identity_into(&dispatch_state, order);
emitter.emit_order_submitted(order);
}
let emitter = emitter.clone();
let ws_client = ws_client.clone();
let http_client = http_client.clone();
get_runtime().spawn(async move {
post_order_batch(
"Bracket child batch",
orders,
requests,
HyperliquidExecGrouping::Na,
builder,
&emitter,
&ws_client,
&http_client,
dispatch_state,
staged_brackets,
clock,
)
.await;
});
}
fn spawn_active_sibling_cancel(
sibling: StagedBracketChild,
emitter: &ExecutionEventEmitter,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
dispatch_state: &WsDispatchState,
) {
let client_order_id = sibling.order.client_order_id();
let Some(cloid) = sibling.request.cloid else {
log::error!("Cannot cancel OUO sibling {client_order_id}: missing CLOID");
return;
};
let venue_order_id = dispatch_state.cached_venue_order_id(&client_order_id);
let action = HyperliquidExecAction::CancelByCloid {
cancels: vec![HyperliquidExecCancelByCloidRequest {
asset: sibling.request.asset,
cloid,
}],
fast: can_fast_cancel_order(Some(sibling.order.order_type())).then_some(true),
};
let emitter = emitter.clone();
let ws_client = ws_client.clone();
let http_client = http_client.clone();
get_runtime().spawn(async move {
match ws_client.post_action_exec(&http_client, &action).await {
Ok(response) if response.is_ok() => {
if let Some(error) = extract_inner_error(&response) {
emitter.emit_order_cancel_rejected(
&sibling.order,
venue_order_id,
&error,
get_atomic_clock_realtime().get_time_ns(),
);
}
}
Ok(response) => {
log::warn!(
"OUO sibling cancel for {client_order_id} returned an ambiguous response; \
awaiting WebSocket or startup reconciliation: {}",
extract_error_message(&response),
);
}
Err(e) => {
log::warn!(
"OUO sibling cancel for {client_order_id} failed; awaiting WebSocket or \
startup reconciliation: {e}",
);
}
}
});
}
fn spawn_active_sibling_resize(
sibling: StagedBracketChild,
target_total_qty: Quantity,
emitter: &ExecutionEventEmitter,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
dispatch_state: &Arc<WsDispatchState>,
) {
let client_order_id = sibling.order.client_order_id();
let Some(old_venue_order_id) = dispatch_state.cached_venue_order_id(&client_order_id) else {
log::warn!(
"Cannot resize OUO sibling {client_order_id}: venue order ID not known; awaiting \
WebSocket or startup reconciliation",
);
return;
};
let filled_qty = dispatch_state
.previous_filled_qty(&client_order_id)
.unwrap_or_else(|| Quantity::zero(target_total_qty.precision));
let Some(order) = build_ouo_resize_request(&sibling, target_total_qty, filled_qty) else {
spawn_active_sibling_cancel(sibling, emitter, ws_client, http_client, dispatch_state);
return;
};
let Some(cloid) = order.cloid else {
log::error!("Cannot resize OUO sibling {client_order_id}: missing CLOID");
return;
};
let generation =
dispatch_state.mark_pending_modify(client_order_id, old_venue_order_id, target_total_qty);
dispatch_state.stash_modify_request(client_order_id, order.clone());
let action = HyperliquidExecAction::Modify {
modify: HyperliquidExecModifyOrderRequest {
oid: HyperliquidExecModifyTarget::Cloid(cloid),
order,
},
};
let ws_client = ws_client.clone();
let http_client = http_client.clone();
let dispatch_state = dispatch_state.clone();
get_runtime().spawn(async move {
match ws_client.post_action_exec(&http_client, &action).await {
Ok(response) if response.is_ok() && extract_inner_error(&response).is_none() => {
log::debug!("OUO sibling resize submitted for {client_order_id}");
}
Ok(response) => {
dispatch_state.clear_modify_generation(&client_order_id, generation);
log::warn!(
"OUO sibling resize for {client_order_id} rejected: {}",
extract_inner_error(&response)
.unwrap_or_else(|| extract_error_message(&response)),
);
}
Err(e) if e.is_transport_error() => {
log::warn!(
"OUO sibling resize transport failure for {client_order_id}: {e}; awaiting \
WebSocket or startup reconciliation",
);
}
Err(e) => {
dispatch_state.clear_modify_generation(&client_order_id, generation);
log::warn!("OUO sibling resize failed for {client_order_id}: {e}");
}
}
});
}
fn build_ouo_resize_request(
sibling: &StagedBracketChild,
target_total_qty: Quantity,
filled_qty: Quantity,
) -> Option<HyperliquidExecPlaceOrderRequest> {
if target_total_qty <= filled_qty {
return None;
}
let mut request = sibling.request.clone();
request.size = (target_total_qty - filled_qty).as_decimal().normalize();
Some(request)
}
struct PostRejectionRoute {
emitter: ExecutionEventEmitter,
ws_client: HyperliquidWebSocketClient,
http_client: HyperliquidHttpClient,
dispatch_state: Arc<WsDispatchState>,
staged_brackets: Arc<Mutex<StagedBracketState>>,
}
impl PostRejectionRoute {
fn new(
emitter: &ExecutionEventEmitter,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
dispatch_state: Arc<WsDispatchState>,
) -> Self {
Self {
emitter: emitter.clone(),
ws_client: ws_client.clone(),
http_client: http_client.clone(),
dispatch_state,
staged_brackets: Arc::new(Mutex::new(StagedBracketState::default())),
}
}
fn with_staged_brackets(
emitter: &ExecutionEventEmitter,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
dispatch_state: Arc<WsDispatchState>,
staged_brackets: Arc<Mutex<StagedBracketState>>,
) -> Self {
Self {
emitter: emitter.clone(),
ws_client: ws_client.clone(),
http_client: http_client.clone(),
dispatch_state,
staged_brackets,
}
}
fn emit_once(
&self,
order: &OrderAny,
reason: &str,
ts_event: UnixNanos,
cloid_hex: &Ustr,
) -> bool {
let client_order_id = order.client_order_id();
let _ = self.dispatch_state.resolve_submission(&client_order_id);
if !self.dispatch_state.insert_filled(client_order_id) {
log::debug!(
"Skipping duplicate post rejection for terminal order {client_order_id}: {reason}",
);
self.ws_client.remove_cloid_mapping(cloid_hex);
self.http_client
.remove_client_order_id_cloid(&client_order_id);
return false;
}
if reason.contains(HYPERLIQUID_BUILDER_FEE_NOT_APPROVED) {
log::warn!(
"Builder fee not approved: complete the one-time 0% builder approval \
(signed by the master wallet). See: {HYPERLIQUID_BUILDER_APPROVAL_DOCS_URL}",
);
}
let normalized_reason = reason.to_lowercase();
let due_post_only = order.is_post_only()
&& (normalized_reason.contains(&HYPERLIQUID_POST_ONLY_WOULD_MATCH.to_lowercase())
|| normalized_reason.contains("post-only order would have immediately matched"));
self.emitter
.emit_order_rejected(order, reason, ts_event, due_post_only);
let active_sibling = self
.staged_brackets
.lock()
.expect(MUTEX_POISONED)
.take_active_sibling(&client_order_id);
if let Some(sibling) = active_sibling {
spawn_active_sibling_cancel(
sibling,
&self.emitter,
&self.ws_client,
&self.http_client,
&self.dispatch_state,
);
}
let staged_children = self
.staged_brackets
.lock()
.expect(MUTEX_POISONED)
.cancel_for_parent(&client_order_id);
for child in staged_children {
self.emitter.emit_order_canceled(&child, None, ts_event);
}
self.dispatch_state
.insert_terminal_cloid(Ustr::from(cloid_hex.as_str()));
self.dispatch_state.cleanup_terminal(&client_order_id);
self.ws_client.remove_cloid_mapping(cloid_hex);
self.http_client
.remove_client_order_id_cloid(&client_order_id);
true
}
fn resolve_without_post_rejection(
&self,
order: &OrderAny,
ts_init: UnixNanos,
cloid_hex: &Ustr,
) {
let client_order_id = order.client_order_id();
let Some(report) = self.dispatch_state.resolve_submission(&client_order_id) else {
return;
};
let is_terminal = !report.order_status.is_open();
let outcome = dispatch_order_event(&report, &self.dispatch_state, &self.emitter, ts_init);
if outcome == DispatchOutcome::External {
self.emitter.send_order_status_report(report);
}
if is_terminal && outcome != DispatchOutcome::Skip {
self.ws_client.remove_cloid_mapping(cloid_hex);
self.http_client
.remove_client_order_id_cloid(&client_order_id);
}
}
}
fn handle_execution_report(
report: ExecutionReport,
dispatch_state: &WsDispatchState,
emitter: &ExecutionEventEmitter,
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
pending_filled_cloids: &mut FifoCache<ClientOrderId, 10_000>,
ts_init: UnixNanos,
) -> Option<(ClientOrderId, u64, HyperliquidExecPlaceOrderRequest)> {
match report {
ExecutionReport::Order(order_report) => {
let is_filled_marker = matches!(order_report.order_status, OrderStatus::Filled);
let is_open = order_report.order_status.is_open();
let client_order_id = order_report.client_order_id;
let outcome = dispatch_order_event(&order_report, dispatch_state, emitter, ts_init);
if outcome == DispatchOutcome::External {
emitter.send_order_status_report(order_report);
}
if let Some(id) = client_order_id
&& !is_open
{
match outcome {
DispatchOutcome::Skip => {}
DispatchOutcome::Tracked if is_filled_marker => {
pending_filled_cloids.add(id);
}
DispatchOutcome::Tracked | DispatchOutcome::External => {
remove_cloid_mapping_for_client_order_id(ws_client, http_client, &id);
}
}
}
client_order_id.and_then(|id| {
dispatch_state
.take_corrective(&id)
.map(|(oid, order)| (id, oid, order))
})
}
ExecutionReport::Fill(fill_report) => {
let client_order_id = fill_report.client_order_id;
let outcome = dispatch_order_fill(&fill_report, dispatch_state, emitter, ts_init);
if outcome == DispatchOutcome::External {
emitter.send_fill_report(fill_report);
}
if let Some(id) = client_order_id
&& pending_filled_cloids.contains(&id)
&& dispatch_state.buffered_fill_count(&id) == 0
{
pending_filled_cloids.remove(&id);
remove_cloid_mapping_for_client_order_id(ws_client, http_client, &id);
}
client_order_id.and_then(|id| {
dispatch_state
.take_corrective(&id)
.map(|(oid, order)| (id, oid, order))
})
}
}
}
fn spawn_corrective_reduce(
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
dispatch_state: &Arc<WsDispatchState>,
client_order_id: ClientOrderId,
oid: u64,
order: HyperliquidExecPlaceOrderRequest,
) {
let ws_client = ws_client.clone();
let http_client = http_client.clone();
let dispatch_state = dispatch_state.clone();
get_runtime().spawn(async move {
let action = HyperliquidExecAction::Modify {
modify: HyperliquidExecModifyOrderRequest {
oid: oid.into(),
order,
},
};
let keep_marker = match ws_client.post_action_exec(&http_client, &action).await {
Ok(resp) if resp.is_ok() && extract_inner_error(&resp).is_none() => {
log::debug!("Corrective reduce acknowledged for {client_order_id} on oid {oid}");
true
}
Ok(resp) => {
let reason =
extract_inner_error(&resp).unwrap_or_else(|| extract_error_message(&resp));
log::warn!(
"Corrective reduce rejected for {client_order_id} on oid {oid}: {reason}"
);
false
}
Err(e) if e.is_transport_error() => {
log::warn!(
"Corrective reduce transport failure for {client_order_id} on oid {oid}: \
{e}; awaiting WS reconciliation",
);
true
}
Err(e) => {
log::warn!("Corrective reduce failed for {client_order_id} on oid {oid}: {e}");
false
}
};
if !keep_marker {
dispatch_state.clear_pending_modify(&client_order_id);
}
});
}
fn remove_cloid_mapping_for_client_order_id(
ws_client: &HyperliquidWebSocketClient,
http_client: &HyperliquidHttpClient,
client_order_id: &ClientOrderId,
) {
let generated_cloid = Cloid::from_client_order_id(*client_order_id);
if let Some(cloid) = http_client.remove_client_order_id_cloid(client_order_id) {
ws_client.remove_cloid_mapping(&Ustr::from(&cloid.to_hex()));
if cloid == generated_cloid {
return;
}
}
ws_client.remove_cloid_mapping(&Ustr::from(&generated_cloid.to_hex()));
}
use crate::common::parse::determine_order_list_grouping;
#[cfg(test)]
mod tests {
use std::sync::Arc;
use nautilus_common::messages::{ExecutionEvent, execution::GenerateOrderStatusReports};
use nautilus_core::{UUID4, UnixNanos, time::get_atomic_clock_realtime};
use nautilus_live::ExecutionEventEmitter;
use nautilus_model::{
enums::{
AccountType, ContingencyType, LiquiditySide, OrderSide, OrderStatus, OrderType,
TimeInForce, TriggerType,
},
events::OrderEventAny,
identifiers::{
AccountId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, VenueOrderId,
},
orders::{Order, OrderAny, limit::LimitOrder, stop_market::StopMarketOrder},
reports::{FillReport, OrderStatusReport},
types::{Currency, Money, Price, Quantity},
};
use nautilus_network::websocket::TransportBackend;
use rstest::rstest;
use rust_decimal::Decimal;
use ustr::Ustr;
use super::{
CancelEntry, ExecutionReport, FifoCache, HyperliquidHttpClient, HyperliquidWebSocketClient,
OrderIdentity, PostRejectionRoute, StagedBracketChild, StagedBracketState, WsDispatchState,
build_ouo_resize_request, can_fast_cancel_order, determine_order_list_grouping,
filter_order_status_reports_for_command, handle_execution_report,
register_order_identity_into, split_fast_cancel_requests, validate_order_for_hyperliquid,
};
use crate::{
common::enums::HyperliquidEnvironment,
http::models::{
Cloid, HyperliquidExecGrouping, HyperliquidExecLimitParams, HyperliquidExecOrderKind,
HyperliquidExecPlaceOrderRequest, HyperliquidExecTif,
},
};
const TEST_INSTRUMENT_ID: &str = "BTC-USD-PERP.HYPERLIQUID";
fn test_emitter() -> (
ExecutionEventEmitter,
tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
) {
let clock = get_atomic_clock_realtime();
let mut emitter = ExecutionEventEmitter::new(
clock,
TraderId::from("TESTER-001"),
AccountId::from("HYPERLIQUID-001"),
AccountType::Margin,
None,
);
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
emitter.set_sender(tx);
(emitter, rx)
}
fn drain_events(
rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
) -> Vec<ExecutionEvent> {
let mut out = Vec::new();
while let Ok(e) = rx.try_recv() {
out.push(e);
}
out
}
fn make_ws_client() -> HyperliquidWebSocketClient {
HyperliquidWebSocketClient::new(
Some("wss://test.invalid".to_string()),
HyperliquidEnvironment::Testnet,
None,
TransportBackend::default(),
None,
)
}
fn make_http_client() -> HyperliquidHttpClient {
HyperliquidHttpClient::new(HyperliquidEnvironment::Testnet, 1, None).unwrap()
}
fn test_identity() -> OrderIdentity {
OrderIdentity {
strategy_id: StrategyId::from("S-001"),
instrument_id: InstrumentId::from(TEST_INSTRUMENT_ID),
order_side: OrderSide::Buy,
order_type: OrderType::Limit,
quantity: Quantity::from("0.0001"),
price: Some(Price::from("56730.0")),
}
}
fn make_status_report(
client_order_id: Option<&str>,
venue_order_id: &str,
status: OrderStatus,
) -> OrderStatusReport {
make_status_report_with_quantity(
client_order_id,
venue_order_id,
status,
Quantity::from("0.0001"),
)
}
fn make_status_report_with_quantity(
client_order_id: Option<&str>,
venue_order_id: &str,
status: OrderStatus,
quantity: Quantity,
) -> OrderStatusReport {
OrderStatusReport::new(
AccountId::from("HYPERLIQUID-001"),
InstrumentId::from(TEST_INSTRUMENT_ID),
client_order_id.map(ClientOrderId::new),
VenueOrderId::new(venue_order_id),
OrderSide::Buy,
OrderType::Limit,
TimeInForce::Gtc,
status,
quantity,
Quantity::from("0"),
UnixNanos::default(),
UnixNanos::default(),
UnixNanos::default(),
Some(UUID4::new()),
)
.with_price(Price::from("56730.0"))
}
fn make_fill_report(
client_order_id: Option<&str>,
venue_order_id: &str,
trade_id: &str,
) -> FillReport {
make_fill_report_with_qty(
client_order_id,
venue_order_id,
trade_id,
Quantity::from("0.0001"),
)
}
fn make_fill_report_with_qty(
client_order_id: Option<&str>,
venue_order_id: &str,
trade_id: &str,
last_qty: Quantity,
) -> FillReport {
FillReport::new(
AccountId::from("HYPERLIQUID-001"),
InstrumentId::from(TEST_INSTRUMENT_ID),
VenueOrderId::new(venue_order_id),
TradeId::new(trade_id),
OrderSide::Buy,
last_qty,
Price::from("56730.0"),
Money::new(0.0, Currency::USD()),
LiquiditySide::Taker,
client_order_id.map(ClientOrderId::new),
None,
UnixNanos::default(),
UnixNanos::default(),
Some(UUID4::new()),
)
}
fn cloid_for(id: &str) -> Ustr {
let cloid = Cloid::from_client_order_id(ClientOrderId::from(id));
Ustr::from(&cloid.to_hex())
}
#[rstest]
fn test_filter_order_status_reports_for_command_filters_open_only() {
let open_report =
make_status_report(Some("O-HER-FILTER-OPEN"), "v-open", OrderStatus::Accepted);
let closed_report =
make_status_report(Some("O-HER-FILTER-CLOSED"), "v-closed", OrderStatus::Filled);
let cmd = order_reports_command(true, None, None);
let filtered =
filter_order_status_reports_for_command(vec![open_report, closed_report], &cmd);
assert_eq!(filtered.len(), 1);
assert_eq!(
filtered[0].client_order_id,
Some(ClientOrderId::from("O-HER-FILTER-OPEN"))
);
}
#[rstest]
fn test_filter_order_status_reports_for_command_filters_time_range_inclusively() {
let mut before = make_status_report(
Some("O-HER-FILTER-BEFORE"),
"v-before",
OrderStatus::Accepted,
);
let mut at_start =
make_status_report(Some("O-HER-FILTER-START"), "v-start", OrderStatus::Accepted);
let mut at_end =
make_status_report(Some("O-HER-FILTER-END"), "v-end", OrderStatus::Accepted);
let mut after =
make_status_report(Some("O-HER-FILTER-AFTER"), "v-after", OrderStatus::Accepted);
before.ts_last = UnixNanos::from(9);
at_start.ts_last = UnixNanos::from(10);
at_end.ts_last = UnixNanos::from(20);
after.ts_last = UnixNanos::from(21);
let cmd =
order_reports_command(false, Some(UnixNanos::from(10)), Some(UnixNanos::from(20)));
let filtered =
filter_order_status_reports_for_command(vec![before, at_start, at_end, after], &cmd);
let filtered_ids: Vec<Option<ClientOrderId>> = filtered
.iter()
.map(|report| report.client_order_id)
.collect();
assert_eq!(
filtered_ids,
vec![
Some(ClientOrderId::from("O-HER-FILTER-START")),
Some(ClientOrderId::from("O-HER-FILTER-END")),
]
);
}
#[rstest]
fn test_filter_order_status_reports_for_command_without_filters_preserves_reports() {
let open_report = make_status_report(
Some("O-HER-FILTER-KEEP-OPEN"),
"v-keep-open",
OrderStatus::Accepted,
);
let closed_report = make_status_report(
Some("O-HER-FILTER-KEEP-CLOSED"),
"v-keep-closed",
OrderStatus::Canceled,
);
let cmd = order_reports_command(false, None, None);
let filtered =
filter_order_status_reports_for_command(vec![open_report, closed_report], &cmd);
let filtered_ids: Vec<Option<ClientOrderId>> = filtered
.iter()
.map(|report| report.client_order_id)
.collect();
assert_eq!(
filtered_ids,
vec![
Some(ClientOrderId::from("O-HER-FILTER-KEEP-OPEN")),
Some(ClientOrderId::from("O-HER-FILTER-KEEP-CLOSED")),
]
);
}
fn order_reports_command(
open_only: bool,
start: Option<UnixNanos>,
end: Option<UnixNanos>,
) -> GenerateOrderStatusReports {
GenerateOrderStatusReports::new(
UUID4::new(),
UnixNanos::default(),
open_only,
None,
start,
end,
None,
None,
)
}
fn limit_order(
id: &str,
reduce_only: bool,
contingency: ContingencyType,
linked_ids: Option<Vec<&str>>,
parent_id: Option<&str>,
) -> OrderAny {
OrderAny::Limit(LimitOrder::new(
TraderId::from("TESTER-001"),
StrategyId::from("S-001"),
InstrumentId::from("ETH-USD-PERP.HYPERLIQUID"),
ClientOrderId::from(id),
OrderSide::Buy,
Quantity::from(1),
Price::from("3000.00"),
TimeInForce::Gtc,
None, false, reduce_only,
false, None, None, None, Some(contingency),
None, linked_ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()),
parent_id.map(ClientOrderId::from),
None, None, None, None, Default::default(),
Default::default(),
))
}
fn stop_order(
id: &str,
reduce_only: bool,
contingency: ContingencyType,
linked_ids: Option<Vec<&str>>,
parent_id: Option<&str>,
) -> OrderAny {
OrderAny::StopMarket(StopMarketOrder::new(
TraderId::from("TESTER-001"),
StrategyId::from("S-001"),
InstrumentId::from("ETH-USD-PERP.HYPERLIQUID"),
ClientOrderId::from(id),
OrderSide::Sell,
Quantity::from(1),
Price::from("2800.00"),
TriggerType::LastPrice,
TimeInForce::Gtc,
None, reduce_only,
false, None, None, None, Some(contingency),
None, linked_ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()),
parent_id.map(ClientOrderId::from),
None, None, None, None, Default::default(),
Default::default(),
))
}
fn staged_child(id: &str, sibling_id: &str) -> StagedBracketChild {
StagedBracketChild {
order: limit_order(
id,
true,
ContingencyType::Ouo,
Some(vec![sibling_id]),
Some("O-PARENT"),
),
request: HyperliquidExecPlaceOrderRequest {
asset: 4,
is_buy: false,
price: Decimal::from(3_000),
size: Decimal::ONE,
reduce_only: true,
kind: HyperliquidExecOrderKind::Limit {
limit: HyperliquidExecLimitParams {
tif: HyperliquidExecTif::Gtc,
},
},
cloid: Some(Cloid::from_client_order_id(ClientOrderId::from(id))),
},
}
}
#[rstest]
fn test_staged_bracket_activation_links_ouo_siblings_once() {
let parent_id = ClientOrderId::from("O-PARENT");
let first_id = ClientOrderId::from("O-CHILD-1");
let second_id = ClientOrderId::from("O-CHILD-2");
let mut state = StagedBracketState::default();
state.stage(
parent_id,
vec![
staged_child(first_id.as_str(), second_id.as_str()),
staged_child(second_id.as_str(), first_id.as_str()),
],
);
let activated = state.activate(&parent_id).expect("staged children");
let sibling = state
.take_active_sibling(&first_id)
.expect("active OUO sibling");
assert_eq!(activated.len(), 2);
assert_eq!(sibling.order.client_order_id(), second_id);
assert!(state.activate(&parent_id).is_none());
assert!(state.take_active_sibling(&second_id).is_none());
}
#[rstest]
fn test_restored_active_bracket_rebuilds_ouo_without_reactivation() {
let parent_id = ClientOrderId::from("O-PARENT");
let first_id = ClientOrderId::from("O-CHILD-1");
let second_id = ClientOrderId::from("O-CHILD-2");
let mut state = StagedBracketState::default();
state.restore_active(&[
staged_child(first_id.as_str(), second_id.as_str()),
staged_child(second_id.as_str(), first_id.as_str()),
]);
let sibling = state
.take_active_sibling(&first_id)
.expect("restored OUO sibling");
assert!(state.activate(&parent_id).is_none());
assert_eq!(sibling.order.client_order_id(), second_id);
assert!(state.take_active_sibling(&second_id).is_none());
}
#[rstest]
fn test_staged_bracket_child_cancel_preserves_other_child_for_parent_fill() {
let parent_id = ClientOrderId::from("O-PARENT");
let first_id = ClientOrderId::from("O-CHILD-1");
let second_id = ClientOrderId::from("O-CHILD-2");
let mut state = StagedBracketState::default();
state.stage(
parent_id,
vec![
staged_child(first_id.as_str(), second_id.as_str()),
staged_child(second_id.as_str(), first_id.as_str()),
],
);
let canceled = state.cancel_child(&first_id).expect("staged child");
let remaining = state.activate(&parent_id).expect("remaining child");
assert_eq!(canceled.client_order_id(), first_id);
assert_eq!(remaining.len(), 1);
assert_eq!(remaining[0].order.client_order_id(), second_id);
}
#[rstest]
fn test_staged_bracket_parent_cancel_returns_all_unsubmitted_children() {
let parent_id = ClientOrderId::from("O-PARENT");
let first_id = ClientOrderId::from("O-CHILD-1");
let second_id = ClientOrderId::from("O-CHILD-2");
let mut state = StagedBracketState::default();
state.stage(
parent_id,
vec![
staged_child(first_id.as_str(), second_id.as_str()),
staged_child(second_id.as_str(), first_id.as_str()),
],
);
let canceled = state.cancel_for_parent(&parent_id);
let canceled_ids = canceled
.iter()
.map(Order::client_order_id)
.collect::<Vec<_>>();
assert_eq!(canceled_ids, vec![first_id, second_id]);
assert!(state.activate(&parent_id).is_none());
}
#[rstest]
fn test_build_ouo_resize_request_sends_sibling_leaves_quantity() {
let sibling = staged_child("O-CHILD-2", "O-CHILD-1");
let request =
build_ouo_resize_request(&sibling, Quantity::from("0.7"), Quantity::from("0.2"))
.expect("resized request");
let exhausted =
build_ouo_resize_request(&sibling, Quantity::from("0.2"), Quantity::from("0.2"));
assert_eq!(request.size, Decimal::new(5, 1));
assert_eq!(request.cloid, sibling.request.cloid);
assert!(exhausted.is_none());
}
#[rstest]
#[case::independent_orders(
vec![
limit_order("O-001", false, ContingencyType::NoContingency, None, None),
limit_order("O-002", false, ContingencyType::NoContingency, None, None),
],
HyperliquidExecGrouping::Na,
)]
#[case::bracket_oto(
vec![
limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002", "O-003"]), None),
limit_order("O-002", true, ContingencyType::Oco, Some(vec!["O-003"]), Some("O-001")),
stop_order("O-003", true, ContingencyType::Oco, Some(vec!["O-002"]), Some("O-001")),
],
HyperliquidExecGrouping::NormalTpsl,
)]
#[case::bracket_oto_with_factory_ouo_children(
vec![
limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002", "O-003"]), None),
limit_order("O-002", true, ContingencyType::Ouo, Some(vec!["O-003"]), Some("O-001")),
stop_order("O-003", true, ContingencyType::Ouo, Some(vec!["O-002"]), Some("O-001")),
],
HyperliquidExecGrouping::NormalTpsl,
)]
#[case::oto_not_bracket_shaped(
vec![
limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002"]), None),
limit_order("O-002", false, ContingencyType::Oto, Some(vec!["O-001"]), None),
],
HyperliquidExecGrouping::Na,
)]
#[case::oco_all_reduce_only(
vec![
limit_order("O-001", true, ContingencyType::Oco, Some(vec!["O-002"]), None),
stop_order("O-002", true, ContingencyType::Oco, Some(vec!["O-001"]), None),
],
HyperliquidExecGrouping::PositionTpsl,
)]
#[case::oco_not_all_reduce_only(
vec![
limit_order("O-001", false, ContingencyType::Oco, Some(vec!["O-002"]), None),
stop_order("O-002", true, ContingencyType::Oco, Some(vec!["O-001"]), None),
],
HyperliquidExecGrouping::Na,
)]
#[case::oto_with_non_oco_children(
vec![
limit_order("O-001", false, ContingencyType::Oto, Some(vec!["O-002", "O-003"]), None),
limit_order("O-002", true, ContingencyType::NoContingency, None, None),
stop_order("O-003", true, ContingencyType::NoContingency, None, None),
],
HyperliquidExecGrouping::Na,
)]
#[case::mixed_oco_and_plain_reduce_only(
vec![
limit_order("O-001", true, ContingencyType::Oco, Some(vec!["O-002"]), None),
stop_order("O-002", true, ContingencyType::NoContingency, None, None),
],
HyperliquidExecGrouping::Na,
)]
#[case::unlinked_oco_reduce_only(
vec![
limit_order("O-001", true, ContingencyType::Oco, Some(vec!["O-099"]), None),
stop_order("O-002", true, ContingencyType::Oco, Some(vec!["O-098"]), None),
],
HyperliquidExecGrouping::Na,
)]
#[case::single_order(
vec![limit_order("O-001", false, ContingencyType::NoContingency, None, None)],
HyperliquidExecGrouping::Na,
)]
fn test_determine_order_list_grouping(
#[case] orders: Vec<OrderAny>,
#[case] expected: HyperliquidExecGrouping,
) {
let result = determine_order_list_grouping(&orders);
assert_eq!(result, expected);
}
#[rstest]
#[case::market(Some(OrderType::Market), true)]
#[case::limit(Some(OrderType::Limit), true)]
#[case::stop_market(Some(OrderType::StopMarket), false)]
#[case::unknown(None, false)]
fn test_can_fast_cancel_order_only_allows_plain_order_types(
#[case] order_type: Option<OrderType>,
#[case] expected: bool,
) {
assert_eq!(can_fast_cancel_order(order_type), expected);
}
#[rstest]
fn test_split_fast_cancel_requests_preserves_request_entry_alignment() {
let requests = vec![
(10_u64, cancel_entry("O-FAST-1", true)),
(20_u64, cancel_entry("O-NORMAL-1", false)),
(30_u64, cancel_entry("O-FAST-2", true)),
(40_u64, cancel_entry("O-NORMAL-2", false)),
];
let (fast_requests, fast_entries, normal_requests, normal_entries) =
split_fast_cancel_requests(requests);
assert_eq!(fast_requests, vec![10, 30]);
assert_eq!(
client_order_ids(&fast_entries),
vec![
ClientOrderId::from("O-FAST-1"),
ClientOrderId::from("O-FAST-2"),
]
);
assert!(fast_entries.iter().all(|entry| entry.fast));
assert_eq!(normal_requests, vec![20, 40]);
assert_eq!(
client_order_ids(&normal_entries),
vec![
ClientOrderId::from("O-NORMAL-1"),
ClientOrderId::from("O-NORMAL-2"),
]
);
assert!(normal_entries.iter().all(|entry| !entry.fast));
}
fn cancel_entry(client_order_id: &str, fast: bool) -> CancelEntry {
CancelEntry {
strategy_id: StrategyId::from("S-001"),
instrument_id: InstrumentId::from(TEST_INSTRUMENT_ID),
client_order_id: ClientOrderId::from(client_order_id),
venue_order_id: Some(VenueOrderId::new("123")),
symbol: Ustr::from("BTC-USD-PERP"),
fast,
}
}
fn client_order_ids(entries: &[CancelEntry]) -> Vec<ClientOrderId> {
entries.iter().map(|entry| entry.client_order_id).collect()
}
fn limit_order_with_flags(id: &str, quote_quantity: bool, post_only: bool) -> OrderAny {
OrderAny::Limit(LimitOrder::new(
TraderId::from("TESTER-001"),
StrategyId::from("S-001"),
InstrumentId::from(TEST_INSTRUMENT_ID),
ClientOrderId::from(id),
OrderSide::Buy,
Quantity::from("0.0001"),
Price::from("56730.0"),
TimeInForce::Gtc,
None,
post_only,
false,
quote_quantity,
None,
None,
None,
Some(ContingencyType::NoContingency),
None,
None,
None,
None,
None,
None,
None,
Default::default(),
Default::default(),
))
}
#[rstest]
fn test_register_order_identity_registers_regular_order() {
let state = WsDispatchState::new();
let order = limit_order_with_flags("O-REG-001", false, false);
register_order_identity_into(&state, &order);
let found = state
.lookup_identity(&ClientOrderId::from("O-REG-001"))
.expect("identity should be registered");
assert_eq!(found.strategy_id, StrategyId::from("S-001"));
assert_eq!(found.instrument_id, InstrumentId::from(TEST_INSTRUMENT_ID));
assert_eq!(found.order_side, OrderSide::Buy);
assert_eq!(found.order_type, OrderType::Limit);
assert_eq!(found.quantity, Quantity::from("0.0001"));
assert_eq!(found.price, Some(Price::from("56730.0")));
}
#[rstest]
fn test_register_order_identity_skips_quote_quantity_order() {
let state = WsDispatchState::new();
let order = limit_order_with_flags("O-QQ-001", true, false);
register_order_identity_into(&state, &order);
assert!(
state
.lookup_identity(&ClientOrderId::from("O-QQ-001"))
.is_none()
);
}
#[rstest]
fn test_handle_execution_report_skip_keeps_cloid_mapping() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-SKIP");
state.register_identity(cid, test_identity());
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("new-voi"));
ws_client.cache_cloid_mapping(cloid_for("O-HER-SKIP"), cid);
let stale_cancel = make_status_report(Some("O-HER-SKIP"), "old-voi", OrderStatus::Canceled);
handle_execution_report(
ExecutionReport::Order(stale_cancel),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
assert!(drain_events(&mut rx).is_empty());
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-SKIP")),
Some(cid)
);
assert!(state.lookup_identity(&cid).is_some());
}
#[rstest]
fn test_handle_execution_report_tracked_terminal_evicts_cloid() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-CANCEL");
state.register_identity(cid, test_identity());
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("v-cancel"));
ws_client.cache_cloid_mapping(cloid_for("O-HER-CANCEL"), cid);
let report = make_status_report(Some("O-HER-CANCEL"), "v-cancel", OrderStatus::Canceled);
handle_execution_report(
ExecutionReport::Order(report),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 1);
assert!(matches!(
events[0],
ExecutionEvent::Order(OrderEventAny::Canceled(_))
));
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-CANCEL")),
None
);
assert!(state.filled_orders.contains(&cid));
}
#[rstest]
fn test_post_rejection_preserves_exact_reason_when_ws_rejection_arrives_first() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = Arc::new(WsDispatchState::new());
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-WS-REJ");
state.register_identity(cid, test_identity());
state.mark_submission_pending(cid);
ws_client.cache_cloid_mapping(cloid_for("O-HER-WS-REJ"), cid);
let report = make_status_report(Some("O-HER-WS-REJ"), "v-rej", OrderStatus::Rejected);
handle_execution_report(
ExecutionReport::Order(report),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
assert!(drain_events(&mut rx).is_empty());
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-WS-REJ")),
Some(cid),
);
let order = limit_order_with_flags("O-HER-WS-REJ", false, true);
let http_client = make_http_client();
let rejection_route =
PostRejectionRoute::new(&emitter, &ws_client, &http_client, state.clone());
let emitted = rejection_route.emit_once(
&order,
"Post only order would have immediately matched, bbo was 56729.0.",
UnixNanos::default(),
&cloid_for("O-HER-WS-REJ"),
);
let events = drain_events(&mut rx);
let ExecutionEvent::Order(OrderEventAny::Rejected(rejected)) = &events[0] else {
panic!("expected OrderRejected, received {:?}", events[0]);
};
assert!(emitted);
assert_eq!(events.len(), 1);
assert_eq!(
rejected.reason.as_str(),
"Post only order would have immediately matched, bbo was 56729.0.",
);
assert!(rejected.due_post_only);
}
#[rstest]
fn test_post_rejection_suppresses_late_raw_cloid_reject() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = Arc::new(WsDispatchState::new());
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-POST-REJ");
let cloid = cloid_for("O-HER-POST-REJ");
let order = limit_order_with_flags("O-HER-POST-REJ", false, true);
state.register_identity(cid, test_identity());
ws_client.cache_cloid_mapping(cloid, cid);
let http_client = make_http_client();
let rejection_route =
PostRejectionRoute::new(&emitter, &ws_client, &http_client, state.clone());
let emitted = rejection_route.emit_once(
&order,
"Post only order would have immediately matched",
UnixNanos::default(),
&cloid,
);
let events = drain_events(&mut rx);
assert!(emitted);
assert_eq!(events.len(), 1);
let ExecutionEvent::Order(OrderEventAny::Rejected(rejected)) = &events[0] else {
panic!("expected OrderRejected, received {:?}", events[0]);
};
assert_eq!(
rejected.reason.as_str(),
"Post only order would have immediately matched",
);
assert!(rejected.due_post_only);
assert_eq!(ws_client.get_cloid_mapping(&cloid), None);
assert!(state.filled_orders.contains(&cid));
assert!(state.terminal_cloid_seen(&cloid));
let late_reject = make_status_report(Some(cloid.as_str()), "v-rej", OrderStatus::Rejected);
handle_execution_report(
ExecutionReport::Order(late_reject),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
assert!(drain_events(&mut rx).is_empty());
}
#[rstest]
fn test_handle_execution_report_filled_marker_then_fill_evicts_on_fill() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-FILL");
state.register_identity(cid, test_identity());
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("v-fill"));
ws_client.cache_cloid_mapping(cloid_for("O-HER-FILL"), cid);
let status_marker = make_status_report(Some("O-HER-FILL"), "v-fill", OrderStatus::Filled);
handle_execution_report(
ExecutionReport::Order(status_marker),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
assert!(drain_events(&mut rx).is_empty());
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-FILL")),
Some(cid)
);
let fill = make_fill_report(Some("O-HER-FILL"), "v-fill", "trade-fill");
handle_execution_report(
ExecutionReport::Fill(fill),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 1);
assert!(matches!(
events[0],
ExecutionEvent::Order(OrderEventAny::Filled(_))
));
assert_eq!(ws_client.get_cloid_mapping(&cloid_for("O-HER-FILL")), None);
}
#[rstest]
fn test_handle_execution_report_fill_under_filled_marker_promotes_and_evicts_cloid() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-BUF");
state.register_identity(cid, test_identity());
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("old-voi"));
state.mark_pending_modify(cid, VenueOrderId::new("old-voi"), test_identity().quantity);
ws_client.cache_cloid_mapping(cloid_for("O-HER-BUF"), cid);
let status_marker = make_status_report(Some("O-HER-BUF"), "new-voi", OrderStatus::Filled);
handle_execution_report(
ExecutionReport::Order(status_marker),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
assert!(pending_cloids.contains(&cid));
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-BUF")),
Some(cid)
);
let fill = make_fill_report(Some("O-HER-BUF"), "new-voi", "trade-buf");
handle_execution_report(
ExecutionReport::Fill(fill),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 2);
assert!(matches!(
events[0],
ExecutionEvent::Order(OrderEventAny::Updated(_))
));
assert!(matches!(
events[1],
ExecutionEvent::Order(OrderEventAny::Filled(_))
));
assert_eq!(state.buffered_fill_count(&cid), 0);
assert!(
!pending_cloids.contains(&cid),
"deferred cleanup must complete once the promoting fill lands",
);
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-BUF")),
None,
"cloid mapping must be evicted after the terminal fill",
);
}
#[rstest]
fn test_cancel_replace_emits_target_total_quantity() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-CR-QTY");
let target_total = Quantity::from("0.00020");
let venue_remaining = Quantity::from("0.00015");
let mut identity = test_identity();
identity.quantity = target_total;
state.register_identity(cid, identity);
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("old-voi"));
state.mark_pending_modify(cid, VenueOrderId::new("old-voi"), target_total);
ws_client.cache_cloid_mapping(cloid_for("O-HER-CR-QTY"), cid);
let accepted = make_status_report_with_quantity(
Some("O-HER-CR-QTY"),
"new-voi",
OrderStatus::Accepted,
venue_remaining,
);
handle_execution_report(
ExecutionReport::Order(accepted),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 1);
match &events[0] {
ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
assert_eq!(
updated.quantity, target_total,
"OrderUpdated must carry the engine's absolute total quantity",
);
assert_eq!(updated.venue_order_id, Some(VenueOrderId::new("new-voi")));
}
other => panic!("expected OrderUpdated, found {other:?}"),
}
let identity = state
.lookup_identity(&cid)
.expect("identity should still be tracked");
assert_eq!(identity.quantity, target_total);
assert!(state.pending_modify(&cid).is_none());
assert!(state.pending_modify_target_qty(&cid).is_none());
assert_eq!(
state.cached_venue_order_id(&cid),
Some(VenueOrderId::new("new-voi")),
);
}
#[rstest]
fn test_cancel_replace_without_marker_falls_back_to_report_quantity() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-CR-EXT");
state.register_identity(cid, test_identity());
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("old-voi"));
ws_client.cache_cloid_mapping(cloid_for("O-HER-CR-EXT"), cid);
let report_qty = Quantity::from("0.0005");
let accepted = make_status_report_with_quantity(
Some("O-HER-CR-EXT"),
"new-voi",
OrderStatus::Accepted,
report_qty,
);
handle_execution_report(
ExecutionReport::Order(accepted),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 1);
match &events[0] {
ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
assert_eq!(updated.quantity, report_qty);
}
other => panic!("expected OrderUpdated, found {other:?}"),
}
}
fn limit_request(size: Decimal) -> HyperliquidExecPlaceOrderRequest {
HyperliquidExecPlaceOrderRequest {
asset: 0,
is_buy: true,
price: "88.949".parse::<Decimal>().unwrap(),
size,
reduce_only: false,
kind: HyperliquidExecOrderKind::Limit {
limit: HyperliquidExecLimitParams {
tif: HyperliquidExecTif::Gtc,
},
},
cloid: None,
}
}
#[rstest]
fn test_cancel_replace_queues_corrective_reduce_on_in_flight_fill() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-4154");
let target_total = Quantity::from("1.000");
let old_voi = "445117664938";
let new_voi = "445117686214";
let mut identity = test_identity();
identity.quantity = target_total;
state.register_identity(cid, identity);
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new(old_voi));
state.mark_pending_modify(cid, VenueOrderId::new(old_voi), target_total);
state.stash_modify_request(cid, limit_request(Decimal::from(1)));
state.record_filled_qty(cid, Quantity::from("0.165"));
let accepted = make_status_report_with_quantity(
Some("O-HER-4154"),
new_voi,
OrderStatus::Accepted,
Quantity::from("0.835"),
);
let corrective = handle_execution_report(
ExecutionReport::Order(accepted),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 1);
match &events[0] {
ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
assert_eq!(updated.quantity, target_total);
assert_eq!(updated.venue_order_id, Some(VenueOrderId::new(new_voi)));
}
other => panic!("expected OrderUpdated, found {other:?}"),
}
let (corr_cid, oid, request) =
corrective.expect("oversized replacement must queue a corrective reduce");
assert_eq!(corr_cid, cid);
assert_eq!(oid, 445_117_686_214);
assert_eq!(request.size, "0.835".parse::<Decimal>().unwrap());
assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new(new_voi)));
assert_eq!(state.pending_modify_target_qty(&cid), Some(target_total));
}
#[rstest]
fn test_cancel_replace_fill_promotion_queues_corrective_reduce() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-FILL-CORR");
let target_total = Quantity::from("1.000");
let old_voi = "445117664938";
let new_voi = "445117686214";
let mut identity = test_identity();
identity.quantity = target_total;
state.register_identity(cid, identity);
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new(old_voi));
state.mark_pending_modify(cid, VenueOrderId::new(old_voi), target_total);
state.stash_modify_request(cid, limit_request(Decimal::from(1)));
state.record_filled_qty(cid, Quantity::from("0.165"));
let fill = make_fill_report_with_qty(
Some("O-HER-FILL-CORR"),
new_voi,
"T-FILL-CORR",
Quantity::from("0.100"),
);
let corrective = handle_execution_report(
ExecutionReport::Fill(fill),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 2);
assert!(matches!(
events[0],
ExecutionEvent::Order(OrderEventAny::Updated(_))
));
assert!(matches!(
events[1],
ExecutionEvent::Order(OrderEventAny::Filled(_))
));
let (corr_cid, oid, request) =
corrective.expect("oversized replacement must queue a corrective reduce");
assert_eq!(corr_cid, cid);
assert_eq!(oid, 445_117_686_214);
assert_eq!(request.size, "0.735".parse::<Decimal>().unwrap());
assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new(new_voi)));
}
#[rstest]
fn test_cancel_replace_no_corrective_without_in_flight_fill() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-4154-NOFILL");
let target_total = Quantity::from("1.000");
let mut identity = test_identity();
identity.quantity = target_total;
state.register_identity(cid, identity);
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("445117664938"));
state.mark_pending_modify(cid, VenueOrderId::new("445117664938"), target_total);
state.stash_modify_request(cid, limit_request(Decimal::from(1)));
let accepted = make_status_report_with_quantity(
Some("O-HER-4154-NOFILL"),
"445117686214",
OrderStatus::Accepted,
target_total,
);
let corrective = handle_execution_report(
ExecutionReport::Order(accepted),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let _ = drain_events(&mut rx);
assert!(corrective.is_none());
assert!(state.pending_modify(&cid).is_none());
assert!(state.take_corrective(&cid).is_none());
assert!(state.modify_request(&cid).is_none());
}
#[rstest]
fn test_cancel_replace_corrective_uses_post_drain_buffered_fill() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-4154-BUF");
let target_total = Quantity::from("1.000");
let new_voi = "445117686214";
let mut identity = test_identity();
identity.quantity = target_total;
state.register_identity(cid, identity);
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("445117664938"));
state.mark_pending_modify(cid, VenueOrderId::new("445117664938"), target_total);
state.stash_modify_request(cid, limit_request(Decimal::from(1)));
let buffered = make_fill_report_with_qty(
Some("O-HER-4154-BUF"),
new_voi,
"trade-buf-4154",
Quantity::from("0.165"),
);
state.buffer_fill(cid, buffered);
let accepted = make_status_report_with_quantity(
Some("O-HER-4154-BUF"),
new_voi,
OrderStatus::Accepted,
Quantity::from("0.835"),
);
let corrective = handle_execution_report(
ExecutionReport::Order(accepted),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let _ = drain_events(&mut rx);
let (_, _, request) =
corrective.expect("buffered fill drained before compute must still queue a corrective");
assert_eq!(request.size, "0.835".parse::<Decimal>().unwrap());
}
#[rstest]
fn test_cancel_replace_no_corrective_when_filled_equals_target() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-4154-EXACT");
let target_total = Quantity::from("1.000");
let mut identity = test_identity();
identity.quantity = target_total;
state.register_identity(cid, identity);
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("445117664938"));
state.mark_pending_modify(cid, VenueOrderId::new("445117664938"), target_total);
state.stash_modify_request(cid, limit_request(Decimal::from(1)));
state.record_filled_qty(cid, target_total);
let accepted = make_status_report_with_quantity(
Some("O-HER-4154-EXACT"),
"445117686214",
OrderStatus::Accepted,
target_total,
);
let corrective = handle_execution_report(
ExecutionReport::Order(accepted),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let _ = drain_events(&mut rx);
assert!(corrective.is_none());
assert!(state.pending_modify(&cid).is_none());
}
#[rstest]
fn test_cancel_replace_chains_second_corrective_reduce() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-4154-CHAIN");
let target_total = Quantity::from("1.000");
let voi3 = "445117699999";
let mut identity = test_identity();
identity.quantity = target_total;
state.register_identity(cid, identity);
state.insert_accepted(cid);
state.record_venue_order_id(cid, VenueOrderId::new("445117686214"));
state.mark_pending_modify(cid, VenueOrderId::new("445117686214"), target_total);
state.stash_modify_request(cid, limit_request("0.835".parse::<Decimal>().unwrap()));
state.record_filled_qty(cid, Quantity::from("0.465"));
let accepted = make_status_report_with_quantity(
Some("O-HER-4154-CHAIN"),
voi3,
OrderStatus::Accepted,
Quantity::from("0.535"),
);
let corrective = handle_execution_report(
ExecutionReport::Order(accepted),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let _ = drain_events(&mut rx);
let (_, oid, request) =
corrective.expect("a further in-flight fill must chain another corrective");
assert_eq!(oid, 445_117_699_999);
assert_eq!(request.size, "0.535".parse::<Decimal>().unwrap());
assert_eq!(state.pending_modify(&cid), Some(VenueOrderId::new(voi3)));
}
#[rstest]
fn test_handle_execution_report_external_terminal_evicts_cloid() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-EXT");
ws_client.cache_cloid_mapping(cloid_for("O-HER-EXT"), cid);
let report = make_status_report(Some("O-HER-EXT"), "v-ext", OrderStatus::Canceled);
handle_execution_report(
ExecutionReport::Order(report),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 1);
assert!(
matches!(events[0], ExecutionEvent::Report(_)),
"external terminal report should forward to the engine as a report",
);
assert_eq!(ws_client.get_cloid_mapping(&cloid_for("O-HER-EXT")), None);
}
#[rstest]
fn test_handle_execution_report_open_status_preserves_cloid() {
let ws_client = make_ws_client();
let (emitter, _rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-OPEN");
state.register_identity(cid, test_identity());
ws_client.cache_cloid_mapping(cloid_for("O-HER-OPEN"), cid);
let report = make_status_report(Some("O-HER-OPEN"), "v-open", OrderStatus::Accepted);
handle_execution_report(
ExecutionReport::Order(report),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-OPEN")),
Some(cid)
);
}
#[rstest]
fn test_handle_execution_report_tracked_accepted_emits_typed_event() {
let ws_client = make_ws_client();
let (emitter, mut rx) = test_emitter();
let state = WsDispatchState::new();
let mut pending_cloids: FifoCache<ClientOrderId, 10_000> = FifoCache::new();
let cid = ClientOrderId::from("O-HER-ACC");
state.register_identity(cid, test_identity());
ws_client.cache_cloid_mapping(cloid_for("O-HER-ACC"), cid);
let report = make_status_report(Some("O-HER-ACC"), "v-acc", OrderStatus::Accepted);
handle_execution_report(
ExecutionReport::Order(report),
&state,
&emitter,
&ws_client,
&make_http_client(),
&mut pending_cloids,
UnixNanos::default(),
);
let events = drain_events(&mut rx);
assert_eq!(events.len(), 1);
assert!(
matches!(events[0], ExecutionEvent::Order(OrderEventAny::Accepted(_))),
"tracked accepted should route through the typed-event path",
);
assert_eq!(
ws_client.get_cloid_mapping(&cloid_for("O-HER-ACC")),
Some(cid)
);
}
fn outcome_limit_order(id: &str, reduce_only: bool) -> OrderAny {
outcome_limit_order_full(id, reduce_only, false, TimeInForce::Gtc)
}
fn outcome_limit_order_full(
id: &str,
reduce_only: bool,
post_only: bool,
time_in_force: TimeInForce,
) -> OrderAny {
OrderAny::Limit(LimitOrder::new(
TraderId::from("TESTER-001"),
StrategyId::from("S-001"),
InstrumentId::from("1-YES-OUTCOME.HYPERLIQUID"),
ClientOrderId::from(id),
OrderSide::Buy,
Quantity::from("1"),
Price::from("0.5000"),
time_in_force,
None,
post_only,
reduce_only,
false,
None,
None,
None,
Some(ContingencyType::NoContingency),
None,
None,
None,
None,
None,
None,
None,
Default::default(),
Default::default(),
))
}
fn outcome_stop_order(id: &str) -> OrderAny {
OrderAny::StopMarket(StopMarketOrder::new(
TraderId::from("TESTER-001"),
StrategyId::from("S-001"),
InstrumentId::from("1-YES-OUTCOME.HYPERLIQUID"),
ClientOrderId::from(id),
OrderSide::Sell,
Quantity::from("1"),
Price::from("0.4000"),
TriggerType::LastPrice,
TimeInForce::Gtc,
None,
false,
false,
None,
None,
None,
Some(ContingencyType::NoContingency),
None,
None,
None,
None,
None,
None,
None,
Default::default(),
Default::default(),
))
}
fn perp_with_unsupported_symbol(id: &str) -> OrderAny {
OrderAny::Limit(LimitOrder::new(
TraderId::from("TESTER-001"),
StrategyId::from("S-001"),
InstrumentId::from("BTC-USD-FOO.HYPERLIQUID"),
ClientOrderId::from(id),
OrderSide::Buy,
Quantity::from("1"),
Price::from("100.0"),
TimeInForce::Gtc,
None,
false,
false,
false,
None,
None,
None,
Some(ContingencyType::NoContingency),
None,
None,
None,
None,
None,
None,
None,
Default::default(),
Default::default(),
))
}
#[rstest]
fn test_validate_accepts_perp_limit_order() {
let order = limit_order(
"O-VAL-PERP",
false,
ContingencyType::NoContingency,
None,
None,
);
validate_order_for_hyperliquid(&order).unwrap();
}
#[rstest]
#[case::gtc_post_only(true, TimeInForce::Gtc)]
#[case::gtc_taker(false, TimeInForce::Gtc)]
#[case::ioc_post_only(true, TimeInForce::Ioc)]
#[case::ioc_taker(false, TimeInForce::Ioc)]
fn test_validate_accepts_outcome_limit_order(
#[case] post_only: bool,
#[case] time_in_force: TimeInForce,
) {
let order = outcome_limit_order_full(
"O-VAL-OUTCOME",
false,
post_only,
time_in_force,
);
validate_order_for_hyperliquid(&order).unwrap();
}
#[rstest]
fn test_validate_rejects_outcome_reduce_only() {
let order = outcome_limit_order("O-VAL-RO", true);
let err = validate_order_for_hyperliquid(&order).unwrap_err();
assert!(
err.to_string().contains("Reduce-only is not supported"),
"unexpected error: {err}",
);
}
#[rstest]
fn test_validate_rejects_outcome_trigger_order() {
let order = outcome_stop_order("O-VAL-TRIG");
let err = validate_order_for_hyperliquid(&order).unwrap_err();
assert!(
err.to_string()
.contains("Trigger order types are not supported"),
"unexpected error: {err}",
);
}
#[rstest]
fn test_validate_rejects_unsupported_symbol_suffix() {
let order = perp_with_unsupported_symbol("O-VAL-BAD");
let err = validate_order_for_hyperliquid(&order).unwrap_err();
assert!(
err.to_string()
.contains("Unsupported instrument symbol format"),
"unexpected error: {err}",
);
}
}