Skip to main content

nautilus_live/execution/
emitter.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live execution event emitter for async event dispatch.
17//!
18//! This module provides [`ExecutionEventEmitter`], which combines event generation (via
19//! [`OrderEventFactory`]) with async dispatch. Adapters use the `emit_*` convenience
20//! methods to generate and send events in a single call.
21//!
22//! # Architecture
23//!
24//! ```text
25//! Adapter
26//! |-- core: ExecutionClientCore    (identity + connection state)
27//! `-- emitter: ExecutionEventEmitter   (event generation + async dispatch)
28//!     |-- factory: OrderEventFactory
29//!     `-- sender: ArcSwapOption<Sender>   (set in start())
30//! ```
31
32use 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/// Event emitter for live trading - combines event generation with async dispatch.
56///
57/// This struct wraps an [`OrderEventFactory`] for event construction and an unbounded
58/// channel sender for async dispatch. It provides `emit_*` convenience methods that
59/// generate and send events in a single call.
60///
61/// The sender is set during the adapter's `start()` phase via [`set_sender`](Self::set_sender).
62/// Clones share the sender slot and observe later sender installations and replacements.
63#[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    /// Creates a new [`ExecutionEventEmitter`] with no sender.
72    ///
73    /// Call [`set_sender`](Self::set_sender) in the adapter's `start()` method.
74    #[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    /// Installs or replaces the sender for this emitter and all its clones.
94    ///
95    /// Call in the adapter's `start()` method.
96    pub fn set_sender(&mut self, sender: tokio::sync::mpsc::UnboundedSender<ExecutionEvent>) {
97        self.sender.store(Some(Arc::new(sender)));
98    }
99
100    /// Returns true if the sender is initialized for this emitter and its clones.
101    #[must_use]
102    pub fn is_initialized(&self) -> bool {
103        self.sender.load().is_some()
104    }
105
106    /// Returns the trader ID.
107    #[must_use]
108    pub fn trader_id(&self) -> TraderId {
109        self.factory.trader_id()
110    }
111
112    /// Returns the account ID.
113    #[must_use]
114    pub fn account_id(&self) -> AccountId {
115        self.factory.account_id()
116    }
117
118    /// Sets the account ID for generated events.
119    pub fn set_account_id(&mut self, account_id: AccountId) {
120        self.factory.set_account_id(account_id);
121    }
122
123    /// Generates and emits an account state event.
124    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    /// Generates and emits an account state event, reporting dispatch failures.
138    ///
139    /// # Errors
140    ///
141    /// Returns an error if the sender is uninitialized or its receiver is closed.
142    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    /// Generates and emits an order denied event.
162    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    /// Generates and emits an order submitted event.
170    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    /// Generates and emits an order rejected event.
176    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    /// Generates and emits an order accepted event.
194    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    /// Generates and emits an order modify rejected event.
207    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    /// Generates and emits an order cancel rejected event.
225    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    /// Generates and emits an order updated event.
243    #[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    /// Generates and emits an order canceled event.
268    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    /// Generates and emits an order triggered event.
281    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    /// Generates and emits an order expired event.
294    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    /// Generates and emits an order filled event.
307    #[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    /// Constructs and emits an order rejected event from raw fields.
338    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    /// Constructs and emits an order modify rejected event from raw fields.
364    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    /// Constructs and emits an order cancel rejected event from raw fields.
390    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    /// Emits an order event.
416    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    /// Emits an order event and returns any channel error to the caller.
423    ///
424    /// # Errors
425    ///
426    /// Returns an error if the sender is uninitialized or its receiver is closed.
427    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    /// Emits a batch of order submitted events as a single channel message.
438    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    /// Emits a batch of order accepted events as a single channel message.
450    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    /// Emits a batch of order canceled events as a single channel message.
462    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    /// Emits an account state event.
474    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    /// Emits an execution report.
491    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    /// Emits an execution report and returns any channel error to the caller.
498    ///
499    /// # Errors
500    ///
501    /// Returns an error if the sender is not initialized or the receiving channel is closed.
502    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    /// Emits an order status report.
513    pub fn send_order_status_report(&self, report: OrderStatusReport) {
514        self.send_execution_report(ExecutionReport::Order(Box::new(report)));
515    }
516
517    /// Emits a fill report.
518    pub fn send_fill_report(&self, report: FillReport) {
519        self.send_execution_report(ExecutionReport::Fill(Box::new(report)));
520    }
521
522    /// Emits an order status report bundled with the fills that produced it.
523    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    /// Emits a position status report.
528    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}