use std::{
collections::VecDeque,
hash::Hash,
sync::{Arc, Mutex},
};
use ahash::AHashMap;
use dashmap::DashMap;
use nautilus_core::{AtomicMap, MUTEX_POISONED, UUID4, UnixNanos, time::AtomicTime};
use nautilus_live::ExecutionEventEmitter;
use nautilus_model::{
enums::{OrderSide, OrderStatus, OrderType},
events::{OrderAccepted, OrderEventAny, OrderFilled, OrderRejected},
identifiers::{
AccountId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, VenueOrderId,
},
instruments::{Instrument, InstrumentAny},
orders::TRIGGERABLE_ORDER_TYPES,
reports::FillReport,
types::{Currency, Money, Quantity},
};
use ustr::Ustr;
use crate::{
common::{
consts::{
OKX_FIELD_CLORDID, OKX_FIELD_SCODE, OKX_FIELD_SMSG, OKX_FIELD_SUBCODE,
OKX_POST_ONLY_CANCEL_REASON, OKX_POST_ONLY_CANCEL_SOURCE, OKX_SUCCESS_CODE,
},
enums::{OKXOrderStatus, OKXOrderType},
parse::{
is_market_price, parse_client_order_id, parse_millisecond_timestamp, parse_price,
parse_quantity,
},
},
http::models::{OKXAccount, OKXCancelAlgoOrderResponse, OKXPosition, OKXSpreadOrder},
websocket::{
client::PendingOrderInfo,
enums::OKXWsOperation,
handler::{is_post_only_auto_cancel, is_unfilled_rpi_cancel},
messages::{ExecutionReport, OKXOrderMsg, OKXWsMessage},
parse::{
OrderStateSnapshot, ParsedOrderEvent, parse_algo_order_msg, parse_order_event,
parse_order_msg, parse_spread_order_event, parse_spread_order_msg,
update_fee_fill_caches,
},
},
};
const DEDUP_CAPACITY: usize = 10_000;
#[derive(Debug)]
pub struct BoundedDedup<K> {
inner: Mutex<BoundedDedupInner<K>>,
capacity: usize,
}
#[derive(Debug)]
struct BoundedDedupInner<K> {
set: AHashMap<K, u64>,
queue: VecDeque<(K, u64)>,
next_seq: u64,
}
impl<K> BoundedDedup<K>
where
K: Eq + Hash + Clone,
{
pub fn new(capacity: usize) -> Self {
Self {
inner: Mutex::new(BoundedDedupInner {
set: AHashMap::with_capacity(capacity),
queue: VecDeque::with_capacity(capacity),
next_seq: 0,
}),
capacity,
}
}
#[allow(clippy::missing_panics_doc, reason = "mutex poisoning is not expected")]
pub fn contains(&self, key: &K) -> bool {
self.inner
.lock()
.expect(MUTEX_POISONED)
.set
.contains_key(key)
}
#[allow(clippy::missing_panics_doc, reason = "mutex poisoning is not expected")]
pub fn check_and_insert(&self, key: K) -> bool {
let mut inner = self.inner.lock().expect(MUTEX_POISONED);
if inner.set.contains_key(&key) {
return true;
}
let seq = inner.next_seq;
inner.next_seq = inner.next_seq.wrapping_add(1);
inner.set.insert(key.clone(), seq);
inner.queue.push_back((key, seq));
while inner.queue.len() > self.capacity
&& let Some((old_key, old_seq)) = inner.queue.pop_front()
{
if inner.set.get(&old_key) == Some(&old_seq) {
inner.set.remove(&old_key);
}
}
false
}
pub fn insert(&self, key: K) {
let _ = self.check_and_insert(key);
}
#[allow(clippy::missing_panics_doc, reason = "mutex poisoning is not expected")]
pub fn remove(&self, key: &K) -> bool {
self.inner
.lock()
.expect(MUTEX_POISONED)
.set
.remove(key)
.is_some()
}
}
#[derive(Debug, Clone)]
pub struct OrderIdentity {
pub instrument_id: InstrumentId,
pub strategy_id: StrategyId,
pub order_side: OrderSide,
pub order_type: OrderType,
}
#[derive(Debug)]
pub struct WsDispatchState {
pub order_identities: DashMap<ClientOrderId, OrderIdentity>,
pub emitted_accepted: BoundedDedup<ClientOrderId>,
pub triggered_orders: BoundedDedup<ClientOrderId>,
pub filled_orders: BoundedDedup<ClientOrderId>,
pub terminal_orders: BoundedDedup<ClientOrderId>,
pub emitted_trades: BoundedDedup<TradeId>,
post_only_rejections: BoundedDedup<Ustr>,
pub(crate) pending_orders: Arc<DashMap<String, PendingOrderInfo>>,
pub(crate) pending_cancels: Arc<DashMap<String, PendingOrderInfo>>,
pub(crate) pending_amends: Arc<DashMap<String, PendingOrderInfo>>,
}
impl Default for WsDispatchState {
fn default() -> Self {
Self {
order_identities: DashMap::new(),
emitted_accepted: BoundedDedup::new(DEDUP_CAPACITY),
triggered_orders: BoundedDedup::new(DEDUP_CAPACITY),
filled_orders: BoundedDedup::new(DEDUP_CAPACITY),
terminal_orders: BoundedDedup::new(DEDUP_CAPACITY),
emitted_trades: BoundedDedup::new(DEDUP_CAPACITY),
post_only_rejections: BoundedDedup::new(DEDUP_CAPACITY),
pending_orders: Arc::new(DashMap::new()),
pending_cancels: Arc::new(DashMap::new()),
pending_amends: Arc::new(DashMap::new()),
}
}
}
impl WsDispatchState {
pub(crate) fn with_pending_maps(
pending_orders: Arc<DashMap<String, PendingOrderInfo>>,
pending_cancels: Arc<DashMap<String, PendingOrderInfo>>,
pending_amends: Arc<DashMap<String, PendingOrderInfo>>,
) -> Self {
Self {
pending_orders,
pending_cancels,
pending_amends,
..Default::default()
}
}
}
impl WsDispatchState {
pub(crate) fn insert_accepted(&self, cid: ClientOrderId) {
self.emitted_accepted.insert(cid);
}
pub(crate) fn insert_triggered(&self, cid: ClientOrderId) {
self.triggered_orders.insert(cid);
}
pub(crate) fn insert_filled(&self, cid: ClientOrderId) {
self.filled_orders.insert(cid);
}
pub(crate) fn insert_terminal(&self, cid: ClientOrderId) {
self.terminal_orders.insert(cid);
}
pub fn check_and_insert_trade(&self, trade_id: TradeId) -> bool {
self.emitted_trades.check_and_insert(trade_id)
}
}
#[expect(clippy::too_many_arguments)]
pub fn dispatch_ws_message(
message: OKXWsMessage,
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
account_id: AccountId,
instruments: &AtomicMap<Ustr, InstrumentAny>,
fee_cache: &mut AHashMap<Ustr, Money>,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
clock: &AtomicTime,
) {
let guard = instruments.load();
let instruments: &AHashMap<Ustr, InstrumentAny> = &guard;
match message {
OKXWsMessage::Orders(order_msgs) => {
let ts_init = clock.get_time_ns();
dispatch_order_messages(
&order_msgs,
emitter,
state,
account_id,
instruments,
fee_cache,
filled_qty_cache,
order_state_cache,
ts_init,
);
}
OKXWsMessage::SpreadOrders(order_msgs) => {
let ts_init = clock.get_time_ns();
dispatch_spread_order_messages(
&order_msgs,
emitter,
state,
account_id,
instruments,
filled_qty_cache,
order_state_cache,
ts_init,
);
}
OKXWsMessage::AlgoOrders(algo_msgs) => {
let ts_init = clock.get_time_ns();
let mut reports = Vec::new();
for msg in algo_msgs {
match parse_algo_order_msg(&msg, account_id, instruments, ts_init) {
Ok(Some(report)) => reports.push(report),
Ok(None) => {}
Err(e) => log::error!("Failed to parse algo order message: {e}"),
}
}
dispatch_execution_reports(reports, emitter, state);
}
OKXWsMessage::Account(data) => {
let ts_init = clock.get_time_ns();
match serde_json::from_value::<Vec<OKXAccount>>(data) {
Ok(accounts) => {
for account in &accounts {
match crate::common::parse::parse_account_state(
account, account_id, ts_init,
) {
Ok(account_state) => emitter.send_account_state(account_state),
Err(e) => log::error!("Failed to parse account state: {e}"),
}
}
}
Err(e) => log::error!("Failed to deserialize account data: {e}"),
}
}
OKXWsMessage::Positions(data) => {
let ts_init = clock.get_time_ns();
match serde_json::from_value::<Vec<OKXPosition>>(data) {
Ok(positions) => {
for position in positions {
let Some(instrument) = instruments.get(&position.inst_id) else {
log::warn!("No cached instrument for position: {}", position.inst_id);
continue;
};
let instrument_id = instrument.id();
let size_precision = instrument.size_precision();
match crate::common::parse::parse_position_status_report(
&position,
account_id,
instrument_id,
size_precision,
ts_init,
) {
Ok(report) => emitter.send_position_report(report),
Err(e) => log::error!("Failed to parse position report: {e}"),
}
}
}
Err(e) => log::error!("Failed to deserialize positions data: {e}"),
}
}
OKXWsMessage::OrderResponse {
id,
op,
code,
msg,
data,
} => {
let ts_init = clock.get_time_ns();
for item in &data {
let s_code = item
.get(OKX_FIELD_SCODE)
.and_then(|v| v.as_str())
.unwrap_or("");
let s_msg = item
.get(OKX_FIELD_SMSG)
.and_then(|v| v.as_str())
.unwrap_or("");
let sub_code = item
.get(OKX_FIELD_SUBCODE)
.and_then(|v| v.as_str())
.unwrap_or("");
let reason = format_order_response_reason(s_code, s_msg, sub_code);
let cl_ord_id = item
.get(OKX_FIELD_CLORDID)
.and_then(|v| v.as_str())
.unwrap_or("");
if s_code == OKX_SUCCESS_CODE {
log::debug!("Order response ok: op={op:?} cl_ord_id={cl_ord_id}");
match op {
OKXWsOperation::Order
| OKXWsOperation::BatchOrders
| OKXWsOperation::OrderAlgo => {
state.pending_orders.remove(cl_ord_id);
}
OKXWsOperation::CancelOrder
| OKXWsOperation::BatchCancelOrders
| OKXWsOperation::MassCancel
| OKXWsOperation::CancelAlgos => {
state.pending_cancels.remove(cl_ord_id);
}
OKXWsOperation::AmendOrder | OKXWsOperation::BatchAmendOrders => {
state.pending_amends.remove(cl_ord_id);
}
_ => {}
}
continue;
}
let Some(client_order_id) = parse_client_order_id(cl_ord_id) else {
log::warn!(
"Order response error without client_order_id: \
op={op:?} s_code={s_code} s_msg={s_msg}"
);
continue;
};
let Some(ident) = state
.order_identities
.get(&client_order_id)
.map(|entry| entry.clone())
else {
log::warn!(
"Order response error for untracked order: \
op={op:?} cl_ord_id={cl_ord_id} s_code={s_code} s_msg={s_msg}"
);
continue;
};
let venue_order_id = item
.get("ordId")
.and_then(|v| v.as_str())
.filter(|s| !s.is_empty())
.map(VenueOrderId::new);
match op {
OKXWsOperation::Order | OKXWsOperation::BatchOrders => {
state.order_identities.remove(&client_order_id);
state.pending_orders.remove(cl_ord_id);
emitter.emit_order_rejected_event(
ident.strategy_id,
ident.instrument_id,
client_order_id,
&reason,
ts_init,
false,
);
}
OKXWsOperation::CancelOrder
| OKXWsOperation::BatchCancelOrders
| OKXWsOperation::MassCancel => {
state.pending_cancels.remove(cl_ord_id);
emitter.emit_order_cancel_rejected_event(
ident.strategy_id,
ident.instrument_id,
client_order_id,
venue_order_id,
&reason,
ts_init,
);
}
OKXWsOperation::AmendOrder | OKXWsOperation::BatchAmendOrders => {
state.pending_amends.remove(cl_ord_id);
emitter.emit_order_modify_rejected_event(
ident.strategy_id,
ident.instrument_id,
client_order_id,
venue_order_id,
&reason,
ts_init,
);
}
_ => {
log::warn!(
"Order response error for unhandled op: \
op={op:?} cl_ord_id={cl_ord_id} s_code={s_code} s_msg={s_msg}"
);
}
}
}
if code != "0" && data.is_empty() {
log::warn!(
"Order response error (no data): id={id:?} op={op:?} code={code} msg={msg}"
);
}
}
OKXWsMessage::SendFailed {
request_id,
client_order_id,
op,
error,
} => {
log::warn!(
"WebSocket send failed without structured venue response: \
request_id={request_id}, client_order_id={client_order_id:?}, \
op={op:?}, awaiting reconciliation: {error}"
);
if let Some(client_order_id) = client_order_id {
let key = client_order_id.as_str();
match op {
Some(
OKXWsOperation::Order
| OKXWsOperation::BatchOrders
| OKXWsOperation::OrderAlgo,
) => {
state.pending_orders.remove(key);
}
Some(
OKXWsOperation::CancelOrder
| OKXWsOperation::BatchCancelOrders
| OKXWsOperation::MassCancel
| OKXWsOperation::CancelAlgos,
) => {
state.pending_cancels.remove(key);
}
Some(OKXWsOperation::AmendOrder | OKXWsOperation::BatchAmendOrders) => {
state.pending_amends.remove(key);
}
_ => {}
}
}
}
OKXWsMessage::ChannelData { channel, .. } => {
log::debug!("Ignoring data channel message on execution client: {channel:?}");
}
OKXWsMessage::BookData { .. }
| OKXWsMessage::RpiBookData { .. }
| OKXWsMessage::Instruments(_) => {
log::debug!("Ignoring data message on execution client");
}
OKXWsMessage::Error(e) => {
log::warn!(
"Websocket error: code={} message={} conn_id={:?}",
e.code,
e.message,
e.conn_id
);
}
OKXWsMessage::Reconnected => {
log::info!("Websocket reconnected");
}
OKXWsMessage::Authenticated => {
log::debug!("Websocket authenticated");
}
}
}
#[expect(clippy::too_many_arguments)]
fn dispatch_order_messages(
order_msgs: &[OKXOrderMsg],
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
account_id: AccountId,
instruments: &AHashMap<Ustr, InstrumentAny>,
fee_cache: &mut AHashMap<Ustr, Money>,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
ts_init: UnixNanos,
) {
for msg in order_msgs {
let Some(instrument) = instruments.get(&msg.inst_id) else {
log::warn!("No instrument for {}, skipping order message", msg.inst_id);
continue;
};
let Some(client_order_id) = parse_client_order_id(&msg.cl_ord_id) else {
log::debug!(
"Order without client_order_id (ord_id={}), sending as report",
msg.ord_id
);
dispatch_order_msg_as_report(
msg,
account_id,
instruments,
fee_cache,
filled_qty_cache,
emitter,
state,
ts_init,
);
continue;
};
let (client_order_id, identity) = match state
.order_identities
.get(&client_order_id)
.map(|r| r.clone())
{
Some(ident) => (client_order_id, Some(ident)),
None => {
if let Some(parent_id) = msg
.algo_cl_ord_id
.as_deref()
.and_then(parse_client_order_id)
{
let parent_ident = state.order_identities.get(&parent_id).map(|r| r.clone());
if parent_ident.is_some() {
(parent_id, parent_ident)
} else {
(client_order_id, None)
}
} else {
(client_order_id, None)
}
}
};
if let Some(ident) = identity {
let is_post_only_cancel = is_post_only_auto_cancel(msg);
if is_post_only_cancel
|| (!state.emitted_accepted.contains(&client_order_id)
&& is_unfilled_rpi_cancel(msg))
{
if is_post_only_cancel {
state.post_only_rejections.insert(msg.ord_id);
}
let ts_event = parse_millisecond_timestamp(msg.u_time);
let reason = if msg.ord_type == OKXOrderType::Rpi {
msg.cancel_source_reason
.as_deref()
.filter(|reason| !reason.is_empty())
.unwrap_or("RPI order canceled before acceptance")
} else {
"Post-only order would have taken liquidity"
};
let rejected = OrderRejected::new(
emitter.trader_id(),
ident.strategy_id,
instrument.id(),
client_order_id,
account_id,
Ustr::from(reason),
UUID4::new(),
ts_event,
ts_init,
false,
true, );
state.order_identities.remove(&client_order_id);
order_state_cache.remove(&client_order_id);
fee_cache.remove(&msg.ord_id);
filled_qty_cache.remove(&msg.ord_id);
emitter.send_order_event(OrderEventAny::Rejected(rejected));
continue;
}
let previous_fee = fee_cache.get(&msg.ord_id).copied();
let previous_filled_qty = filled_qty_cache.get(&msg.ord_id).copied();
let previous_state = order_state_cache.get(&client_order_id);
match parse_order_event(
msg,
client_order_id,
account_id,
emitter.trader_id(),
ident.strategy_id,
instrument,
previous_fee,
previous_filled_qty,
previous_state,
ts_init,
) {
Ok(event) => {
update_order_caches(
msg,
instrument,
client_order_id,
fee_cache,
filled_qty_cache,
order_state_cache,
);
dispatch_parsed_order_event(
event,
client_order_id,
account_id,
VenueOrderId::new(msg.ord_id),
&ident,
instrument,
msg.state,
emitter,
state,
order_state_cache,
ts_init,
);
}
Err(e) => log::error!("Failed to parse order event for {client_order_id}: {e}"),
}
} else if is_post_only_auto_cancel(msg) && state.post_only_rejections.contains(&msg.ord_id)
{
log::debug!(
"Skipping replayed post-only rejection for {client_order_id}: ord_id={}",
msg.ord_id
);
} else {
log::debug!(
"Untracked order {client_order_id} (ord_id={}), sending as report for reconciliation",
msg.ord_id
);
dispatch_order_msg_as_report(
msg,
account_id,
instruments,
fee_cache,
filled_qty_cache,
emitter,
state,
ts_init,
);
}
}
}
#[expect(clippy::too_many_arguments)]
fn dispatch_spread_order_messages(
order_msgs: &[OKXSpreadOrder],
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
account_id: AccountId,
instruments: &AHashMap<Ustr, InstrumentAny>,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
ts_init: UnixNanos,
) {
for msg in order_msgs {
let Some(instrument) = instruments.get(&msg.sprd_id) else {
log::warn!(
"No instrument for {}, skipping spread order message",
msg.sprd_id
);
continue;
};
let Some(client_order_id) = parse_client_order_id(msg.cl_ord_id.as_str()) else {
log::debug!(
"Spread order without client_order_id (ord_id={}), sending as report",
msg.ord_id
);
dispatch_spread_order_msg_as_report(
msg,
account_id,
instruments,
filled_qty_cache,
emitter,
state,
ts_init,
);
continue;
};
let identity = state
.order_identities
.get(&client_order_id)
.map(|r| r.clone());
if let Some(ident) = identity {
if is_spread_post_only_auto_cancel(msg) {
let ts_event = msg
.u_time
.or(msg.c_time)
.map_or(ts_init, parse_millisecond_timestamp);
let rejected = OrderRejected::new(
emitter.trader_id(),
ident.strategy_id,
instrument.id(),
client_order_id,
account_id,
Ustr::from(OKX_POST_ONLY_CANCEL_REASON),
UUID4::new(),
ts_event,
ts_init,
false,
true,
);
state.order_identities.remove(&client_order_id);
order_state_cache.remove(&client_order_id);
filled_qty_cache.remove(&msg.ord_id);
emitter.send_order_event(OrderEventAny::Rejected(rejected));
continue;
}
let previous_filled_qty = filled_qty_cache.get(&msg.ord_id).copied();
let previous_state = order_state_cache.get(&client_order_id);
match parse_spread_order_event(
msg,
client_order_id,
account_id,
emitter.trader_id(),
ident.strategy_id,
instrument,
previous_filled_qty,
previous_state,
ts_init,
) {
Ok(event) => {
update_spread_order_caches(
msg,
instrument,
client_order_id,
filled_qty_cache,
order_state_cache,
);
dispatch_parsed_order_event(
event,
client_order_id,
account_id,
VenueOrderId::new(msg.ord_id.as_str()),
&ident,
instrument,
msg.state,
emitter,
state,
order_state_cache,
ts_init,
);
}
Err(e) => {
log::error!("Failed to parse spread order event for {client_order_id}: {e}");
}
}
} else {
log::debug!(
"Untracked spread order {client_order_id} (ord_id={}), sending as report for reconciliation",
msg.ord_id
);
dispatch_spread_order_msg_as_report(
msg,
account_id,
instruments,
filled_qty_cache,
emitter,
state,
ts_init,
);
}
}
}
#[expect(clippy::too_many_arguments)]
fn dispatch_parsed_order_event(
event: ParsedOrderEvent,
client_order_id: ClientOrderId,
account_id: AccountId,
venue_order_id: VenueOrderId,
identity: &OrderIdentity,
instrument: &InstrumentAny,
venue_status: OKXOrderStatus,
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
ts_init: UnixNanos,
) {
let is_terminal;
match event {
ParsedOrderEvent::Accepted(e) => {
if state.emitted_accepted.contains(&client_order_id)
|| state.filled_orders.contains(&client_order_id)
|| state.triggered_orders.contains(&client_order_id)
|| state.terminal_orders.contains(&client_order_id)
{
log::debug!("Skipping duplicate Accepted for {client_order_id}");
return;
}
state.insert_accepted(client_order_id);
is_terminal = false;
emitter.send_order_event(OrderEventAny::Accepted(e));
}
ParsedOrderEvent::Triggered(e) => {
if state.filled_orders.contains(&client_order_id) {
log::debug!("Skipping stale Triggered for {client_order_id} (already filled)");
return;
}
if !TRIGGERABLE_ORDER_TYPES.contains(&identity.order_type) {
log::debug!(
"Skipping OrderTriggered for {} order {client_order_id}: market-style stops have no TRIGGERED state",
identity.order_type,
);
state.insert_triggered(client_order_id);
return;
}
ensure_accepted_emitted(
client_order_id,
account_id,
venue_order_id,
identity,
emitter,
state,
ts_init,
);
state.insert_triggered(client_order_id);
is_terminal = false;
emitter.send_order_event(OrderEventAny::Triggered(e));
}
ParsedOrderEvent::Canceled(e) => {
ensure_accepted_emitted(
client_order_id,
account_id,
venue_order_id,
identity,
emitter,
state,
ts_init,
);
state.triggered_orders.remove(&client_order_id);
state.filled_orders.remove(&client_order_id);
is_terminal = true;
emitter.send_order_event(OrderEventAny::Canceled(e));
}
ParsedOrderEvent::Expired(e) => {
ensure_accepted_emitted(
client_order_id,
account_id,
venue_order_id,
identity,
emitter,
state,
ts_init,
);
state.triggered_orders.remove(&client_order_id);
state.filled_orders.remove(&client_order_id);
is_terminal = true;
emitter.send_order_event(OrderEventAny::Expired(e));
}
ParsedOrderEvent::Updated(e) => {
ensure_accepted_emitted(
client_order_id,
account_id,
venue_order_id,
identity,
emitter,
state,
ts_init,
);
is_terminal = false;
emitter.send_order_event(OrderEventAny::Updated(e));
}
ParsedOrderEvent::Fill(fill_report) => {
let is_duplicate = state.check_and_insert_trade(fill_report.trade_id);
is_terminal = venue_status == OKXOrderStatus::Filled;
if is_duplicate {
log::debug!(
"Skipping duplicate fill for {client_order_id}: trade_id={}",
fill_report.trade_id
);
} else {
ensure_accepted_emitted(
client_order_id,
account_id,
venue_order_id,
identity,
emitter,
state,
ts_init,
);
state.insert_filled(client_order_id);
state.triggered_orders.remove(&client_order_id);
let filled = fill_report_to_order_filled(
&fill_report,
emitter.trader_id(),
identity,
instrument.quote_currency(),
);
emitter.send_order_event(OrderEventAny::Filled(filled));
}
}
ParsedOrderEvent::StatusOnly(report) => {
is_terminal = matches!(
report.order_status,
OrderStatus::Filled | OrderStatus::Canceled | OrderStatus::Expired
);
emitter.send_order_status_report(*report);
}
ParsedOrderEvent::Skipped => return,
}
if is_terminal {
state.insert_terminal(client_order_id);
state.order_identities.remove(&client_order_id);
state.emitted_accepted.remove(&client_order_id);
order_state_cache.remove(&client_order_id);
}
}
fn ensure_accepted_emitted(
client_order_id: ClientOrderId,
account_id: AccountId,
venue_order_id: VenueOrderId,
identity: &OrderIdentity,
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
ts_init: UnixNanos,
) {
if state.emitted_accepted.contains(&client_order_id) {
return;
}
state.insert_accepted(client_order_id);
let accepted = OrderAccepted::new(
emitter.trader_id(),
identity.strategy_id,
identity.instrument_id,
client_order_id,
venue_order_id,
account_id,
UUID4::new(),
ts_init,
ts_init,
false,
);
emitter.send_order_event(OrderEventAny::Accepted(accepted));
}
fn fill_report_to_order_filled(
report: &FillReport,
trader_id: TraderId,
identity: &OrderIdentity,
quote_currency: Currency,
) -> OrderFilled {
OrderFilled::new(
trader_id,
identity.strategy_id,
report.instrument_id,
report
.client_order_id
.expect("tracked order has client_order_id"),
report.venue_order_id,
report.account_id,
report.trade_id,
identity.order_side,
identity.order_type,
report.last_qty,
report.last_px,
quote_currency,
report.liquidity_side,
UUID4::new(),
report.ts_event,
report.ts_init,
false,
report.venue_position_id,
Some(report.commission),
None,
)
}
#[expect(clippy::too_many_arguments)]
fn dispatch_order_msg_as_report(
msg: &OKXOrderMsg,
account_id: AccountId,
instruments: &AHashMap<Ustr, InstrumentAny>,
fee_cache: &mut AHashMap<Ustr, Money>,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
ts_init: UnixNanos,
) {
match parse_order_msg(
msg,
account_id,
instruments,
fee_cache,
filled_qty_cache,
ts_init,
) {
Ok(report) => {
if let Some(instrument) = instruments.get(&msg.inst_id) {
update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
}
dispatch_execution_reports(vec![report], emitter, state);
}
Err(e) => log::error!("Failed to parse order message as report: {e}"),
}
}
fn dispatch_spread_order_msg_as_report(
msg: &OKXSpreadOrder,
account_id: AccountId,
instruments: &AHashMap<Ustr, InstrumentAny>,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
ts_init: UnixNanos,
) {
match parse_spread_order_msg(msg, account_id, instruments, filled_qty_cache, ts_init) {
Ok(report) => {
if let Some(instrument) = instruments.get(&msg.sprd_id) {
update_spread_fill_cache(msg, instrument, filled_qty_cache);
}
dispatch_execution_reports(vec![report], emitter, state);
}
Err(e) => log::error!("Failed to parse spread order message as report: {e}"),
}
}
fn update_order_caches(
msg: &OKXOrderMsg,
instrument: &InstrumentAny,
client_order_id: ClientOrderId,
fee_cache: &mut AHashMap<Ustr, Money>,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
) {
update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
let venue_order_id = VenueOrderId::new(msg.ord_id);
let quantity = parse_quantity(&msg.sz, instrument.size_precision()).unwrap_or_default();
let price = if is_market_price(&msg.px) {
None
} else {
parse_price(&msg.px, instrument.price_precision()).ok()
};
order_state_cache.insert(
client_order_id,
OrderStateSnapshot {
venue_order_id,
quantity,
price,
},
);
}
fn update_spread_order_caches(
msg: &OKXSpreadOrder,
instrument: &InstrumentAny,
client_order_id: ClientOrderId,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
) {
update_spread_fill_cache(msg, instrument, filled_qty_cache);
let venue_order_id = VenueOrderId::new(msg.ord_id.as_str());
let quantity = parse_quantity(&msg.sz, instrument.size_precision()).unwrap_or_default();
let price = if is_market_price(&msg.px) {
None
} else {
parse_price(&msg.px, instrument.price_precision()).ok()
};
order_state_cache.insert(
client_order_id,
OrderStateSnapshot {
venue_order_id,
quantity,
price,
},
);
}
fn update_spread_fill_cache(
msg: &OKXSpreadOrder,
instrument: &InstrumentAny,
filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
) {
if !msg.acc_fill_sz.is_empty()
&& msg.acc_fill_sz != "0"
&& let Ok(qty) = parse_quantity(&msg.acc_fill_sz, instrument.size_precision())
{
filled_qty_cache.insert(msg.ord_id, qty);
}
}
fn is_spread_post_only_auto_cancel(msg: &OKXSpreadOrder) -> bool {
msg.state == OKXOrderStatus::Canceled && msg.cancel_source == OKX_POST_ONLY_CANCEL_SOURCE
}
pub fn dispatch_execution_reports(
reports: Vec<ExecutionReport>,
emitter: &ExecutionEventEmitter,
state: &WsDispatchState,
) {
log::debug!("Processing {} execution report(s)", reports.len());
for report in reports {
match report {
ExecutionReport::Order(order_report) => {
if let Some(cid) = order_report.client_order_id {
match order_report.order_status {
#[allow(clippy::collapsible_match)]
OrderStatus::Accepted => {
if state.terminal_orders.contains(&cid)
|| state.filled_orders.contains(&cid)
|| state.triggered_orders.contains(&cid)
{
log::debug!(
"Skipping stale OrderStatusReport(Accepted) \
for {cid} (order already terminal)"
);
continue;
}
}
OrderStatus::Triggered => {
if state.filled_orders.contains(&cid) {
log::debug!(
"Skipping stale OrderStatusReport(Triggered) \
for {cid} (already filled)"
);
continue;
}
state.insert_triggered(cid);
}
OrderStatus::Filled => {
state.insert_filled(cid);
state.insert_terminal(cid);
state.triggered_orders.remove(&cid);
}
OrderStatus::Canceled | OrderStatus::Expired | OrderStatus::Rejected => {
state.insert_terminal(cid);
state.triggered_orders.remove(&cid);
state.filled_orders.remove(&cid);
}
_ => {}
}
}
emitter.send_order_status_report(order_report);
}
ExecutionReport::Fill(fill_report) => {
if state.check_and_insert_trade(fill_report.trade_id) {
log::debug!(
"Skipping duplicate fill report: trade_id={}",
fill_report.trade_id
);
continue;
}
if let Some(cid) = fill_report.client_order_id {
state.insert_filled(cid);
state.triggered_orders.remove(&cid);
}
emitter.send_fill_report(fill_report);
}
}
}
}
fn format_order_response_reason(s_code: &str, s_msg: &str, sub_code: &str) -> String {
match (s_msg.is_empty(), sub_code.is_empty(), s_code.is_empty()) {
(false, true, _) => s_msg.to_string(),
(false, false, _) => format!("{s_msg} (subCode={sub_code})"),
(true, false, false) => format!("sCode={s_code} subCode={sub_code}"),
(true, false, true) => format!("subCode={sub_code}"),
(true, true, false) => format!("sCode={s_code}"),
(true, true, true) => String::new(),
}
}
#[derive(Debug, Clone)]
pub struct AlgoCancelContext {
pub client_order_id: ClientOrderId,
pub instrument_id: InstrumentId,
pub strategy_id: StrategyId,
pub venue_order_id: Option<VenueOrderId>,
}
pub fn emit_algo_cancel_rejections(
responses: &[OKXCancelAlgoOrderResponse],
contexts: &[AlgoCancelContext],
emitter: &ExecutionEventEmitter,
clock: &'static AtomicTime,
) {
for (i, item) in responses.iter().enumerate() {
let code = item.s_code.as_deref().unwrap_or(OKX_SUCCESS_CODE);
if code == OKX_SUCCESS_CODE {
continue;
}
let msg = item.s_msg.as_deref().unwrap_or("");
if let Some(ctx) = contexts.get(i) {
let ts = clock.get_time_ns();
emitter.emit_order_cancel_rejected_event(
ctx.strategy_id,
ctx.instrument_id,
ctx.client_order_id,
ctx.venue_order_id,
msg,
ts,
);
} else {
log::warn!(
"Algo cancel rejected but no context at index {i}: \
algo_id={} sCode={code} sMsg={msg}",
item.algo_id
);
}
}
}
pub fn emit_batch_cancel_failure(
contexts: &[AlgoCancelContext],
error: &str,
_emitter: &ExecutionEventEmitter,
_clock: &'static AtomicTime,
) {
for ctx in contexts {
log::warn!(
"Ambiguous algo batch cancel failure for {}, awaiting reconciliation: {error}",
ctx.client_order_id
);
}
}
#[cfg(test)]
mod tests {
use rstest::rstest;
use super::{BoundedDedup, format_order_response_reason};
#[rstest]
#[case("51000", "Rejected", "", "Rejected")]
#[case("51000", "Rejected", "51004", "Rejected (subCode=51004)")]
#[case("51000", "", "51004", "sCode=51000 subCode=51004")]
#[case("51000", "", "", "sCode=51000")]
#[case("", "", "51004", "subCode=51004")]
#[case("", "", "", "")]
fn test_format_order_response_reason(
#[case] s_code: &str,
#[case] s_msg: &str,
#[case] sub_code: &str,
#[case] expected: &str,
) {
assert_eq!(
format_order_response_reason(s_code, s_msg, sub_code),
expected
);
}
#[rstest]
fn test_bounded_dedup_check_and_insert_returns_false_on_first_insert() {
let dedup = BoundedDedup::<u32>::new(4);
assert!(!dedup.check_and_insert(1));
assert!(dedup.contains(&1));
}
#[rstest]
fn test_bounded_dedup_check_and_insert_returns_true_on_duplicate() {
let dedup = BoundedDedup::<u32>::new(4);
dedup.insert(1);
assert!(dedup.check_and_insert(1));
}
#[rstest]
fn test_bounded_dedup_evicts_oldest_on_overflow() {
let dedup = BoundedDedup::<u32>::new(3);
dedup.insert(1);
dedup.insert(2);
dedup.insert(3);
dedup.insert(4);
assert!(!dedup.contains(&1));
assert!(dedup.contains(&2));
assert!(dedup.contains(&3));
assert!(dedup.contains(&4));
}
#[rstest]
fn test_bounded_dedup_evicted_key_is_not_treated_as_duplicate() {
let dedup = BoundedDedup::<u32>::new(2);
dedup.insert(1);
dedup.insert(2);
dedup.insert(3);
assert!(!dedup.check_and_insert(1));
assert!(dedup.contains(&1));
}
#[rstest]
fn test_bounded_dedup_remove_drops_entry() {
let dedup = BoundedDedup::<u32>::new(4);
dedup.insert(1);
assert!(dedup.remove(&1));
assert!(!dedup.contains(&1));
assert!(!dedup.remove(&1));
}
#[rstest]
fn test_bounded_dedup_remove_then_reinsert_does_not_double_count() {
let dedup = BoundedDedup::<u32>::new(2);
for k in 0u32..1000 {
dedup.insert(k);
dedup.remove(&k);
}
assert!(!dedup.contains(&0));
assert!(!dedup.contains(&500));
}
#[rstest]
fn test_bounded_dedup_reinsert_survives_stale_marker_eviction() {
let dedup = BoundedDedup::<u32>::new(3);
dedup.insert(1);
dedup.remove(&1);
dedup.insert(2);
dedup.insert(3);
dedup.insert(1);
assert!(dedup.contains(&1));
assert!(dedup.contains(&2));
assert!(dedup.contains(&3));
}
}