1use std::sync::Arc;
33
34use arc_swap::ArcSwapOption;
35use nautilus_common::{
36 factories::OrderEventFactory,
37 messages::{ExecutionEvent, ExecutionReport},
38};
39use nautilus_core::{Params, UUID4, UnixNanos, time::AtomicTime};
40use nautilus_model::{
41 enums::{AccountType, LiquiditySide},
42 events::{
43 AccountState, OrderAcceptedBatch, OrderCancelRejected, OrderCanceledBatch, OrderEventAny,
44 OrderModifyRejected, OrderRejected, OrderSubmittedBatch,
45 },
46 identifiers::{
47 AccountId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId, TraderId,
48 VenueOrderId,
49 },
50 orders::OrderAny,
51 reports::{FillReport, OrderStatusReport, PositionStatusReport},
52 types::{AccountBalance, Currency, MarginBalance, Money, Price, Quantity},
53};
54
55#[derive(Debug, Clone)]
64pub struct ExecutionEventEmitter {
65 clock: &'static AtomicTime,
66 factory: OrderEventFactory,
67 sender: Arc<ArcSwapOption<tokio::sync::mpsc::UnboundedSender<ExecutionEvent>>>,
68}
69
70impl ExecutionEventEmitter {
71 #[must_use]
75 pub fn new(
76 clock: &'static AtomicTime,
77 trader_id: TraderId,
78 account_id: AccountId,
79 account_type: AccountType,
80 base_currency: Option<Currency>,
81 ) -> Self {
82 Self {
83 clock,
84 factory: OrderEventFactory::new(trader_id, account_id, account_type, base_currency),
85 sender: Arc::new(ArcSwapOption::empty()),
86 }
87 }
88
89 fn ts_init(&self) -> UnixNanos {
90 self.clock.get_time_ns()
91 }
92
93 pub fn set_sender(&mut self, sender: tokio::sync::mpsc::UnboundedSender<ExecutionEvent>) {
97 self.sender.store(Some(Arc::new(sender)));
98 }
99
100 #[must_use]
102 pub fn is_initialized(&self) -> bool {
103 self.sender.load().is_some()
104 }
105
106 #[must_use]
108 pub fn trader_id(&self) -> TraderId {
109 self.factory.trader_id()
110 }
111
112 #[must_use]
114 pub fn account_id(&self) -> AccountId {
115 self.factory.account_id()
116 }
117
118 pub fn set_account_id(&mut self, account_id: AccountId) {
120 self.factory.set_account_id(account_id);
121 }
122
123 pub fn emit_account_state(
125 &self,
126 balances: Vec<AccountBalance>,
127 margins: Vec<MarginBalance>,
128 reported: bool,
129 ts_event: UnixNanos,
130 info: Option<Params>,
131 ) {
132 if let Err(e) = self.try_emit_account_state(balances, margins, reported, ts_event, info) {
133 log::warn!("{e}");
134 }
135 }
136
137 pub fn try_emit_account_state(
143 &self,
144 balances: Vec<AccountBalance>,
145 margins: Vec<MarginBalance>,
146 reported: bool,
147 ts_event: UnixNanos,
148 info: Option<Params>,
149 ) -> anyhow::Result<()> {
150 let state = self.factory.generate_account_state(
151 balances,
152 margins,
153 reported,
154 ts_event,
155 self.ts_init(),
156 info,
157 );
158 self.try_send_account_state(state)
159 }
160
161 pub fn emit_order_denied(&self, order: &OrderAny, reason: &str) {
163 let event = self
164 .factory
165 .generate_order_denied(order, reason, self.ts_init());
166 self.send_order_event(event);
167 }
168
169 pub fn emit_order_submitted(&self, order: &OrderAny) {
171 let event = self.factory.generate_order_submitted(order, self.ts_init());
172 self.send_order_event(event);
173 }
174
175 pub fn emit_order_rejected(
177 &self,
178 order: &OrderAny,
179 reason: &str,
180 ts_event: UnixNanos,
181 due_post_only: bool,
182 ) {
183 let event = self.factory.generate_order_rejected(
184 order,
185 reason,
186 ts_event,
187 self.ts_init(),
188 due_post_only,
189 );
190 self.send_order_event(event);
191 }
192
193 pub fn emit_order_accepted(
195 &self,
196 order: &OrderAny,
197 venue_order_id: VenueOrderId,
198 ts_event: UnixNanos,
199 ) {
200 let event =
201 self.factory
202 .generate_order_accepted(order, venue_order_id, ts_event, self.ts_init());
203 self.send_order_event(event);
204 }
205
206 pub fn emit_order_modify_rejected(
208 &self,
209 order: &OrderAny,
210 venue_order_id: Option<VenueOrderId>,
211 reason: &str,
212 ts_event: UnixNanos,
213 ) {
214 let event = self.factory.generate_order_modify_rejected(
215 order,
216 venue_order_id,
217 reason,
218 ts_event,
219 self.ts_init(),
220 );
221 self.send_order_event(event);
222 }
223
224 pub fn emit_order_cancel_rejected(
226 &self,
227 order: &OrderAny,
228 venue_order_id: Option<VenueOrderId>,
229 reason: &str,
230 ts_event: UnixNanos,
231 ) {
232 let event = self.factory.generate_order_cancel_rejected(
233 order,
234 venue_order_id,
235 reason,
236 ts_event,
237 self.ts_init(),
238 );
239 self.send_order_event(event);
240 }
241
242 #[expect(clippy::too_many_arguments)]
244 pub fn emit_order_updated(
245 &self,
246 order: &OrderAny,
247 venue_order_id: VenueOrderId,
248 quantity: Quantity,
249 price: Option<Price>,
250 trigger_price: Option<Price>,
251 protection_price: Option<Price>,
252 ts_event: UnixNanos,
253 ) {
254 let event = self.factory.generate_order_updated(
255 order,
256 venue_order_id,
257 quantity,
258 price,
259 trigger_price,
260 protection_price,
261 ts_event,
262 self.ts_init(),
263 );
264 self.send_order_event(event);
265 }
266
267 pub fn emit_order_canceled(
269 &self,
270 order: &OrderAny,
271 venue_order_id: Option<VenueOrderId>,
272 ts_event: UnixNanos,
273 ) {
274 let event =
275 self.factory
276 .generate_order_canceled(order, venue_order_id, ts_event, self.ts_init());
277 self.send_order_event(event);
278 }
279
280 pub fn emit_order_triggered(
282 &self,
283 order: &OrderAny,
284 venue_order_id: Option<VenueOrderId>,
285 ts_event: UnixNanos,
286 ) {
287 let event =
288 self.factory
289 .generate_order_triggered(order, venue_order_id, ts_event, self.ts_init());
290 self.send_order_event(event);
291 }
292
293 pub fn emit_order_expired(
295 &self,
296 order: &OrderAny,
297 venue_order_id: Option<VenueOrderId>,
298 ts_event: UnixNanos,
299 ) {
300 let event =
301 self.factory
302 .generate_order_expired(order, venue_order_id, ts_event, self.ts_init());
303 self.send_order_event(event);
304 }
305
306 #[expect(clippy::too_many_arguments)]
308 pub fn emit_order_filled(
309 &self,
310 order: &OrderAny,
311 venue_order_id: VenueOrderId,
312 venue_position_id: Option<PositionId>,
313 trade_id: TradeId,
314 last_qty: Quantity,
315 last_px: Price,
316 quote_currency: Currency,
317 commission: Option<Money>,
318 liquidity_side: LiquiditySide,
319 ts_event: UnixNanos,
320 ) {
321 let event = self.factory.generate_order_filled(
322 order,
323 venue_order_id,
324 venue_position_id,
325 trade_id,
326 last_qty,
327 last_px,
328 quote_currency,
329 commission,
330 liquidity_side,
331 ts_event,
332 self.ts_init(),
333 );
334 self.send_order_event(event);
335 }
336
337 pub fn emit_order_rejected_event(
339 &self,
340 strategy_id: StrategyId,
341 instrument_id: InstrumentId,
342 client_order_id: ClientOrderId,
343 reason: &str,
344 ts_event: UnixNanos,
345 due_post_only: bool,
346 ) {
347 let event = OrderRejected::new(
348 self.factory.trader_id(),
349 strategy_id,
350 instrument_id,
351 client_order_id,
352 self.factory.account_id(),
353 reason.into(),
354 UUID4::new(),
355 ts_event,
356 self.ts_init(),
357 false,
358 due_post_only,
359 );
360 self.send_order_event(OrderEventAny::Rejected(event));
361 }
362
363 pub fn emit_order_modify_rejected_event(
365 &self,
366 strategy_id: StrategyId,
367 instrument_id: InstrumentId,
368 client_order_id: ClientOrderId,
369 venue_order_id: Option<VenueOrderId>,
370 reason: &str,
371 ts_event: UnixNanos,
372 ) {
373 let event = OrderModifyRejected::new(
374 self.factory.trader_id(),
375 strategy_id,
376 instrument_id,
377 client_order_id,
378 reason.into(),
379 UUID4::new(),
380 ts_event,
381 self.ts_init(),
382 false,
383 venue_order_id,
384 Some(self.factory.account_id()),
385 );
386 self.send_order_event(OrderEventAny::ModifyRejected(event));
387 }
388
389 pub fn emit_order_cancel_rejected_event(
391 &self,
392 strategy_id: StrategyId,
393 instrument_id: InstrumentId,
394 client_order_id: ClientOrderId,
395 venue_order_id: Option<VenueOrderId>,
396 reason: &str,
397 ts_event: UnixNanos,
398 ) {
399 let event = OrderCancelRejected::new(
400 self.factory.trader_id(),
401 strategy_id,
402 instrument_id,
403 client_order_id,
404 reason.into(),
405 UUID4::new(),
406 ts_event,
407 self.ts_init(),
408 false,
409 venue_order_id,
410 Some(self.factory.account_id()),
411 );
412 self.send_order_event(OrderEventAny::CancelRejected(event));
413 }
414
415 pub fn send_order_event(&self, event: OrderEventAny) {
417 if let Err(e) = self.try_send_order_event(event) {
418 log::warn!("{e}");
419 }
420 }
421
422 pub fn try_send_order_event(&self, event: OrderEventAny) -> anyhow::Result<()> {
428 let sender = self.sender.load();
429 let sender = sender
430 .as_ref()
431 .ok_or_else(|| anyhow::anyhow!("Cannot send order event: sender not initialized"))?;
432 sender
433 .send(ExecutionEvent::Order(event))
434 .map_err(|e| anyhow::anyhow!("Failed to send order event: {e}"))
435 }
436
437 pub fn send_order_submitted_batch(&self, batch: OrderSubmittedBatch) {
439 let sender = self.sender.load();
440 if let Some(sender) = sender.as_ref() {
441 if let Err(e) = sender.send(ExecutionEvent::OrderSubmittedBatch(batch)) {
442 log::warn!("Failed to send order submitted batch: {e}");
443 }
444 } else {
445 log::warn!("Cannot send order submitted batch: sender not initialized");
446 }
447 }
448
449 pub fn send_order_accepted_batch(&self, batch: OrderAcceptedBatch) {
451 let sender = self.sender.load();
452 if let Some(sender) = sender.as_ref() {
453 if let Err(e) = sender.send(ExecutionEvent::OrderAcceptedBatch(batch)) {
454 log::warn!("Failed to send order accepted batch: {e}");
455 }
456 } else {
457 log::warn!("Cannot send order accepted batch: sender not initialized");
458 }
459 }
460
461 pub fn send_order_canceled_batch(&self, batch: OrderCanceledBatch) {
463 let sender = self.sender.load();
464 if let Some(sender) = sender.as_ref() {
465 if let Err(e) = sender.send(ExecutionEvent::OrderCanceledBatch(batch)) {
466 log::warn!("Failed to send order canceled batch: {e}");
467 }
468 } else {
469 log::warn!("Cannot send order canceled batch: sender not initialized");
470 }
471 }
472
473 pub fn send_account_state(&self, state: AccountState) {
475 if let Err(e) = self.try_send_account_state(state) {
476 log::warn!("{e}");
477 }
478 }
479
480 fn try_send_account_state(&self, state: AccountState) -> anyhow::Result<()> {
481 let sender = self.sender.load();
482 let sender = sender
483 .as_ref()
484 .ok_or_else(|| anyhow::anyhow!("Cannot send account state: sender not initialized"))?;
485 sender
486 .send(ExecutionEvent::Account(state))
487 .map_err(|e| anyhow::anyhow!("Failed to send account state: {e}"))
488 }
489
490 pub fn send_execution_report(&self, report: ExecutionReport) {
492 if let Err(e) = self.try_send_execution_report(report) {
493 log::warn!("{e}");
494 }
495 }
496
497 pub fn try_send_execution_report(&self, report: ExecutionReport) -> anyhow::Result<()> {
503 let sender = self.sender.load();
504 let sender = sender.as_ref().ok_or_else(|| {
505 anyhow::anyhow!("Cannot send execution report: sender not initialized")
506 })?;
507 sender
508 .send(ExecutionEvent::Report(report))
509 .map_err(|e| anyhow::anyhow!("Failed to send execution report: {e}"))
510 }
511
512 pub fn send_order_status_report(&self, report: OrderStatusReport) {
514 self.send_execution_report(ExecutionReport::Order(Box::new(report)));
515 }
516
517 pub fn send_fill_report(&self, report: FillReport) {
519 self.send_execution_report(ExecutionReport::Fill(Box::new(report)));
520 }
521
522 pub fn send_order_with_fills(&self, report: OrderStatusReport, fills: Vec<FillReport>) {
524 self.send_execution_report(ExecutionReport::OrderWithFills(Box::new(report), fills));
525 }
526
527 pub fn send_position_report(&self, report: PositionStatusReport) {
529 self.send_execution_report(ExecutionReport::Position(Box::new(report)));
530 }
531}
532
533#[cfg(test)]
534mod tests {
535 use nautilus_core::time::get_atomic_clock_static;
536 use nautilus_model::events::order::spec::OrderSubmittedSpec;
537 use rstest::rstest;
538
539 use super::*;
540
541 fn create_emitter() -> ExecutionEventEmitter {
542 ExecutionEventEmitter::new(
543 get_atomic_clock_static(),
544 TraderId::from("TRADER-001"),
545 AccountId::from("SIM-001"),
546 AccountType::Cash,
547 None,
548 )
549 }
550
551 fn create_order_event() -> OrderEventAny {
552 OrderEventAny::Submitted(
553 OrderSubmittedSpec::builder()
554 .client_order_id(ClientOrderId::from("O-001"))
555 .build(),
556 )
557 }
558
559 #[rstest]
560 fn test_clone_before_set_sender_observes_sender() {
561 let mut emitter = create_emitter();
562 let cloned = emitter.clone();
563 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
564
565 emitter.set_sender(tx);
566
567 assert!(cloned.is_initialized());
568 cloned.send_order_event(create_order_event());
569 assert!(matches!(
570 rx.try_recv(),
571 Ok(ExecutionEvent::Order(OrderEventAny::Submitted(_)))
572 ));
573 }
574
575 #[rstest]
576 fn test_set_sender_replaces_existing_sender() {
577 let mut emitter = create_emitter();
578 let (tx_a, mut rx_a) = tokio::sync::mpsc::unbounded_channel();
579 let (tx_b, mut rx_b) = tokio::sync::mpsc::unbounded_channel();
580
581 emitter.set_sender(tx_a);
582 emitter.set_sender(tx_b);
583 emitter.send_order_event(create_order_event());
584
585 assert!(matches!(
586 rx_a.try_recv(),
587 Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
588 ));
589 assert!(matches!(
590 rx_b.try_recv(),
591 Ok(ExecutionEvent::Order(OrderEventAny::Submitted(_)))
592 ));
593 }
594
595 #[rstest]
596 fn test_never_initialized_sender_drops_or_errors() {
597 let emitter = create_emitter();
598
599 emitter.send_order_event(create_order_event());
600 let error = emitter
601 .try_send_order_event(create_order_event())
602 .unwrap_err();
603
604 assert_eq!(
605 error.to_string(),
606 "Cannot send order event: sender not initialized"
607 );
608 }
609}