nautilus_binance/common/
dispatch.rs1use dashmap::DashMap;
25use nautilus_common::cache::fifo::FifoCache;
26use nautilus_core::{UUID4, UnixNanos};
27use nautilus_live::ExecutionEventEmitter;
28use nautilus_model::{
29 enums::{OrderSide, OrderType},
30 events::{OrderAccepted, OrderEventAny},
31 identifiers::{AccountId, ClientOrderId, InstrumentId, PositionId, StrategyId, VenueOrderId},
32 types::{Price, Quantity},
33};
34use parking_lot::Mutex;
35
36#[derive(Debug, Clone, Copy)]
38pub enum PendingOperation {
39 Place,
40 Cancel,
41 Modify,
42}
43
44#[derive(Debug, Clone)]
50pub struct PendingRequest {
51 pub client_order_id: ClientOrderId,
52 pub venue_order_id: Option<VenueOrderId>,
53 pub operation: PendingOperation,
54}
55
56#[derive(Debug, Clone)]
61pub struct OrderIdentity {
62 pub instrument_id: InstrumentId,
63 pub strategy_id: StrategyId,
64 pub order_side: OrderSide,
65 pub order_type: OrderType,
66 pub price: Option<Price>,
67 pub quantity: Quantity,
68 pub venue_position_id: Option<PositionId>,
69}
70
71#[derive(Debug, Clone, Copy)]
72struct AlgoOrderIds {
73 algo: VenueOrderId,
74 current: VenueOrderId,
75}
76
77#[derive(Debug)]
83pub struct WsDispatchState {
84 pub order_identities: DashMap<ClientOrderId, OrderIdentity>,
85 pub pending_requests: DashMap<String, PendingRequest>,
86 algo_order_ids: DashMap<ClientOrderId, AlgoOrderIds>,
87 emitted_accepted: Mutex<FifoCache<ClientOrderId, 10_000>>,
88 filled_orders: Mutex<FifoCache<ClientOrderId, 10_000>>,
89}
90
91impl Default for WsDispatchState {
92 fn default() -> Self {
93 Self {
94 order_identities: DashMap::new(),
95 pending_requests: DashMap::new(),
96 algo_order_ids: DashMap::new(),
97 emitted_accepted: Mutex::new(FifoCache::new()),
98 filled_orders: Mutex::new(FifoCache::new()),
99 }
100 }
101}
102
103impl WsDispatchState {
104 pub fn has_emitted_accepted(&self, cid: &ClientOrderId) -> bool {
105 self.emitted_accepted.lock().contains(cid)
106 }
107
108 pub fn insert_accepted(&self, cid: ClientOrderId) {
110 self.emitted_accepted.lock().add(cid);
111 }
112
113 pub fn has_filled(&self, cid: &ClientOrderId) -> bool {
114 self.filled_orders.lock().contains(cid)
115 }
116
117 pub fn insert_filled(&self, cid: ClientOrderId) {
119 self.filled_orders.lock().add(cid);
120 }
121
122 pub fn insert_algo_order_id(&self, cid: ClientOrderId, venue_order_id: VenueOrderId) {
123 self.algo_order_ids.entry(cid).or_insert(AlgoOrderIds {
124 algo: venue_order_id,
125 current: venue_order_id,
126 });
127 }
128
129 pub fn promote_algo_order_id(
134 &self,
135 cid: ClientOrderId,
136 venue_order_id: VenueOrderId,
137 ) -> Option<bool> {
138 let mut ids = self.algo_order_ids.get_mut(&cid)?;
139 let changed = ids.current != venue_order_id;
140 ids.current = venue_order_id;
141 Some(changed)
142 }
143
144 pub fn promoted_algo_order_id(&self, cid: &ClientOrderId) -> Option<VenueOrderId> {
146 self.algo_order_ids
147 .get(cid)
148 .and_then(|ids| (ids.current != ids.algo).then_some(ids.current))
149 }
150
151 pub fn cleanup_terminal(&self, cid: ClientOrderId) {
153 self.order_identities.remove(&cid);
154 self.algo_order_ids.remove(&cid);
155 self.emitted_accepted.lock().remove(&cid);
156 self.filled_orders.lock().remove(&cid);
157 }
158}
159
160pub fn ensure_accepted_emitted(
164 client_order_id: ClientOrderId,
165 account_id: AccountId,
166 venue_order_id: VenueOrderId,
167 identity: &OrderIdentity,
168 emitter: &ExecutionEventEmitter,
169 state: &WsDispatchState,
170 ts_init: UnixNanos,
171) {
172 if state.has_emitted_accepted(&client_order_id) {
173 return;
174 }
175 state.insert_accepted(client_order_id);
176 let accepted = OrderAccepted::new(
177 emitter.trader_id(),
178 identity.strategy_id,
179 identity.instrument_id,
180 client_order_id,
181 venue_order_id,
182 account_id,
183 UUID4::new(),
184 ts_init,
185 ts_init,
186 false,
187 );
188 emitter.send_order_event(OrderEventAny::Accepted(accepted));
189}