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>   (shared slot, installed at create() and start())
30//! ```
31
32use std::sync::Arc;
33
34use arc_swap::ArcSwapOption;
35use nautilus_common::{
36    factories::OrderEventFactory,
37    live::sender::EventSender,
38    messages::{ExecutionEvent, ExecutionReport},
39};
40use nautilus_core::{Params, UUID4, UnixNanos, time::AtomicTime};
41use nautilus_model::{
42    enums::{AccountType, LiquiditySide},
43    events::{
44        AccountState, OrderAcceptedBatch, OrderCancelRejected, OrderCanceledBatch, OrderEventAny,
45        OrderModifyRejected, OrderRejected, OrderSubmittedBatch,
46    },
47    identifiers::{
48        AccountId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId, TraderId,
49        VenueOrderId,
50    },
51    orders::OrderAny,
52    reports::{FillReport, OrderStatusReport, PositionStatusReport},
53    types::{AccountBalance, Currency, MarginBalance, Money, Price, Quantity},
54};
55
56/// Event emitter for live trading - combines event generation with async dispatch.
57///
58/// This struct wraps an [`OrderEventFactory`] for event construction and an unbounded
59/// channel sender for async dispatch. It provides `emit_*` convenience methods that
60/// generate and send events in a single call.
61///
62/// The sender is installed via [`set_sender`](Self::set_sender) in the client's `start`, and in the
63/// execution client factory's `create` when the calling thread has one.
64/// Clones share the sender slot and observe later sender installations and replacements.
65#[derive(Debug, Clone)]
66pub struct ExecutionEventEmitter {
67    clock: &'static AtomicTime,
68    factory: OrderEventFactory,
69    sender: Arc<ArcSwapOption<EventSender<ExecutionEvent>>>,
70}
71
72impl ExecutionEventEmitter {
73    /// Creates a new [`ExecutionEventEmitter`] with no sender.
74    ///
75    /// Call [`set_sender`](Self::set_sender) in the client's `start`, and in the factory's `create`
76    /// when `try_get_exec_event_sender` returns `Some`.
77    #[must_use]
78    pub fn new(
79        clock: &'static AtomicTime,
80        trader_id: TraderId,
81        account_id: AccountId,
82        account_type: AccountType,
83        base_currency: Option<Currency>,
84    ) -> Self {
85        Self {
86            clock,
87            factory: OrderEventFactory::new(trader_id, account_id, account_type, base_currency),
88            sender: Arc::new(ArcSwapOption::empty()),
89        }
90    }
91
92    fn ts_init(&self) -> UnixNanos {
93        self.clock.get_time_ns()
94    }
95
96    /// Installs or replaces the sender for this emitter and all its clones.
97    ///
98    /// The slot is shared, so a clone taken before this call observes the sender it installs.
99    /// Events emitted before any install are dropped with a warning.
100    ///
101    /// Call in the client's `start`, resolved from `get_exec_event_sender`: `LiveNode` rebinds the
102    /// runner's senders on the calling thread before it starts clients, so the `start` install is
103    /// the authoritative one. Call it in the execution client factory's `create` as well when
104    /// `try_get_exec_event_sender` returns `Some`; `None` there means the calling thread has no
105    /// bound senders and is not a construction failure. See the adapter guide for hosts that drive
106    /// a client outside `LiveNode`.
107    pub fn set_sender(&mut self, sender: impl Into<EventSender<ExecutionEvent>>) {
108        self.sender.store(Some(Arc::new(sender.into())));
109    }
110
111    /// Returns true if the sender is initialized for this emitter and its clones.
112    #[must_use]
113    pub fn is_initialized(&self) -> bool {
114        self.sender.load().is_some()
115    }
116
117    /// Returns the trader ID.
118    #[must_use]
119    pub fn trader_id(&self) -> TraderId {
120        self.factory.trader_id()
121    }
122
123    /// Returns the account ID.
124    #[must_use]
125    pub fn account_id(&self) -> AccountId {
126        self.factory.account_id()
127    }
128
129    /// Sets the account ID for generated events.
130    pub fn set_account_id(&mut self, account_id: AccountId) {
131        self.factory.set_account_id(account_id);
132    }
133
134    /// Generates and emits an account state event.
135    pub fn emit_account_state(
136        &self,
137        balances: Vec<AccountBalance>,
138        margins: Vec<MarginBalance>,
139        reported: bool,
140        ts_event: UnixNanos,
141        info: Option<Params>,
142    ) {
143        if let Err(e) = self.try_emit_account_state(balances, margins, reported, ts_event, info) {
144            log::warn!("{e}");
145        }
146    }
147
148    /// Generates and emits an account state event, reporting dispatch failures.
149    ///
150    /// # Errors
151    ///
152    /// Returns an error if the sender is uninitialized or its receiver is closed.
153    pub fn try_emit_account_state(
154        &self,
155        balances: Vec<AccountBalance>,
156        margins: Vec<MarginBalance>,
157        reported: bool,
158        ts_event: UnixNanos,
159        info: Option<Params>,
160    ) -> anyhow::Result<()> {
161        let state = self.factory.generate_account_state(
162            balances,
163            margins,
164            reported,
165            ts_event,
166            self.ts_init(),
167            info,
168        );
169        self.try_send_account_state(state)
170    }
171
172    /// Generates and emits an order denied event.
173    pub fn emit_order_denied(&self, order: &OrderAny, reason: &str) {
174        let event = self
175            .factory
176            .generate_order_denied(order, reason, self.ts_init());
177        self.send_order_event(event);
178    }
179
180    /// Generates and emits an order submitted event.
181    pub fn emit_order_submitted(&self, order: &OrderAny) {
182        let event = self.factory.generate_order_submitted(order, self.ts_init());
183        self.send_order_event(event);
184    }
185
186    /// Generates and emits an order rejected event.
187    pub fn emit_order_rejected(
188        &self,
189        order: &OrderAny,
190        reason: &str,
191        ts_event: UnixNanos,
192        due_post_only: bool,
193    ) {
194        let event = self.factory.generate_order_rejected(
195            order,
196            reason,
197            ts_event,
198            self.ts_init(),
199            due_post_only,
200        );
201        self.send_order_event(event);
202    }
203
204    /// Generates and emits an order accepted event.
205    pub fn emit_order_accepted(
206        &self,
207        order: &OrderAny,
208        venue_order_id: VenueOrderId,
209        ts_event: UnixNanos,
210    ) {
211        let event =
212            self.factory
213                .generate_order_accepted(order, venue_order_id, ts_event, self.ts_init());
214        self.send_order_event(event);
215    }
216
217    /// Generates and emits an order modify rejected event.
218    pub fn emit_order_modify_rejected(
219        &self,
220        order: &OrderAny,
221        venue_order_id: Option<VenueOrderId>,
222        reason: &str,
223        ts_event: UnixNanos,
224    ) {
225        let event = self.factory.generate_order_modify_rejected(
226            order,
227            venue_order_id,
228            reason,
229            ts_event,
230            self.ts_init(),
231        );
232        self.send_order_event(event);
233    }
234
235    /// Generates and emits an order cancel rejected event.
236    pub fn emit_order_cancel_rejected(
237        &self,
238        order: &OrderAny,
239        venue_order_id: Option<VenueOrderId>,
240        reason: &str,
241        ts_event: UnixNanos,
242    ) {
243        let event = self.factory.generate_order_cancel_rejected(
244            order,
245            venue_order_id,
246            reason,
247            ts_event,
248            self.ts_init(),
249        );
250        self.send_order_event(event);
251    }
252
253    /// Generates and emits an order updated event.
254    #[expect(clippy::too_many_arguments)]
255    pub fn emit_order_updated(
256        &self,
257        order: &OrderAny,
258        venue_order_id: VenueOrderId,
259        quantity: Quantity,
260        price: Option<Price>,
261        trigger_price: Option<Price>,
262        protection_price: Option<Price>,
263        ts_event: UnixNanos,
264    ) {
265        let event = self.factory.generate_order_updated(
266            order,
267            venue_order_id,
268            quantity,
269            price,
270            trigger_price,
271            protection_price,
272            ts_event,
273            self.ts_init(),
274        );
275        self.send_order_event(event);
276    }
277
278    /// Generates and emits an order canceled event.
279    pub fn emit_order_canceled(
280        &self,
281        order: &OrderAny,
282        venue_order_id: Option<VenueOrderId>,
283        ts_event: UnixNanos,
284    ) {
285        let event =
286            self.factory
287                .generate_order_canceled(order, venue_order_id, ts_event, self.ts_init());
288        self.send_order_event(event);
289    }
290
291    /// Generates and emits an order triggered event.
292    pub fn emit_order_triggered(
293        &self,
294        order: &OrderAny,
295        venue_order_id: Option<VenueOrderId>,
296        ts_event: UnixNanos,
297    ) {
298        let event =
299            self.factory
300                .generate_order_triggered(order, venue_order_id, ts_event, self.ts_init());
301        self.send_order_event(event);
302    }
303
304    /// Generates and emits an order expired event.
305    pub fn emit_order_expired(
306        &self,
307        order: &OrderAny,
308        venue_order_id: Option<VenueOrderId>,
309        ts_event: UnixNanos,
310    ) {
311        let event =
312            self.factory
313                .generate_order_expired(order, venue_order_id, ts_event, self.ts_init());
314        self.send_order_event(event);
315    }
316
317    /// Generates and emits an order filled event.
318    #[expect(clippy::too_many_arguments)]
319    pub fn emit_order_filled(
320        &self,
321        order: &OrderAny,
322        venue_order_id: VenueOrderId,
323        venue_position_id: Option<PositionId>,
324        trade_id: TradeId,
325        last_qty: Quantity,
326        last_px: Price,
327        quote_currency: Currency,
328        commission: Option<Money>,
329        liquidity_side: LiquiditySide,
330        ts_event: UnixNanos,
331    ) {
332        let event = self.factory.generate_order_filled(
333            order,
334            venue_order_id,
335            venue_position_id,
336            trade_id,
337            last_qty,
338            last_px,
339            quote_currency,
340            commission,
341            liquidity_side,
342            ts_event,
343            self.ts_init(),
344        );
345        self.send_order_event(event);
346    }
347
348    /// Constructs and emits an order rejected event from raw fields.
349    pub fn emit_order_rejected_event(
350        &self,
351        strategy_id: StrategyId,
352        instrument_id: InstrumentId,
353        client_order_id: ClientOrderId,
354        reason: &str,
355        ts_event: UnixNanos,
356        due_post_only: bool,
357    ) {
358        let event = OrderRejected::new(
359            self.factory.trader_id(),
360            strategy_id,
361            instrument_id,
362            client_order_id,
363            self.factory.account_id(),
364            reason.into(),
365            UUID4::new(),
366            ts_event,
367            self.ts_init(),
368            false,
369            due_post_only,
370        );
371        self.send_order_event(OrderEventAny::Rejected(event));
372    }
373
374    /// Constructs and emits an order modify rejected event from raw fields.
375    pub fn emit_order_modify_rejected_event(
376        &self,
377        strategy_id: StrategyId,
378        instrument_id: InstrumentId,
379        client_order_id: ClientOrderId,
380        venue_order_id: Option<VenueOrderId>,
381        reason: &str,
382        ts_event: UnixNanos,
383    ) {
384        let event = OrderModifyRejected::new(
385            self.factory.trader_id(),
386            strategy_id,
387            instrument_id,
388            client_order_id,
389            reason.into(),
390            UUID4::new(),
391            ts_event,
392            self.ts_init(),
393            false,
394            venue_order_id,
395            Some(self.factory.account_id()),
396        );
397        self.send_order_event(OrderEventAny::ModifyRejected(event));
398    }
399
400    /// Constructs and emits an order cancel rejected event from raw fields.
401    pub fn emit_order_cancel_rejected_event(
402        &self,
403        strategy_id: StrategyId,
404        instrument_id: InstrumentId,
405        client_order_id: ClientOrderId,
406        venue_order_id: Option<VenueOrderId>,
407        reason: &str,
408        ts_event: UnixNanos,
409    ) {
410        let event = OrderCancelRejected::new(
411            self.factory.trader_id(),
412            strategy_id,
413            instrument_id,
414            client_order_id,
415            reason.into(),
416            UUID4::new(),
417            ts_event,
418            self.ts_init(),
419            false,
420            venue_order_id,
421            Some(self.factory.account_id()),
422        );
423        self.send_order_event(OrderEventAny::CancelRejected(event));
424    }
425
426    /// Emits an order event.
427    pub fn send_order_event(&self, event: OrderEventAny) {
428        if let Err(e) = self.try_send_order_event(event) {
429            log::warn!("{e}");
430        }
431    }
432
433    /// Emits an order event and returns any channel error to the caller.
434    ///
435    /// # Errors
436    ///
437    /// Returns an error if the sender is uninitialized or its receiver is closed.
438    pub fn try_send_order_event(&self, event: OrderEventAny) -> anyhow::Result<()> {
439        let sender = self.sender.load();
440        let sender = sender
441            .as_ref()
442            .ok_or_else(|| anyhow::anyhow!("Cannot send order event: sender not initialized"))?;
443        sender
444            .send(ExecutionEvent::Order(event))
445            .map_err(|e| anyhow::anyhow!("Failed to send order event: {e}"))
446    }
447
448    /// Emits a batch of order submitted events as a single channel message.
449    pub fn send_order_submitted_batch(&self, batch: OrderSubmittedBatch) {
450        let sender = self.sender.load();
451        if let Some(sender) = sender.as_ref() {
452            if let Err(e) = sender.send(ExecutionEvent::OrderSubmittedBatch(batch)) {
453                log::warn!("Failed to send order submitted batch: {e}");
454            }
455        } else {
456            log::warn!("Cannot send order submitted batch: sender not initialized");
457        }
458    }
459
460    /// Emits a batch of order accepted events as a single channel message.
461    pub fn send_order_accepted_batch(&self, batch: OrderAcceptedBatch) {
462        let sender = self.sender.load();
463        if let Some(sender) = sender.as_ref() {
464            if let Err(e) = sender.send(ExecutionEvent::OrderAcceptedBatch(batch)) {
465                log::warn!("Failed to send order accepted batch: {e}");
466            }
467        } else {
468            log::warn!("Cannot send order accepted batch: sender not initialized");
469        }
470    }
471
472    /// Emits a batch of order canceled events as a single channel message.
473    pub fn send_order_canceled_batch(&self, batch: OrderCanceledBatch) {
474        let sender = self.sender.load();
475        if let Some(sender) = sender.as_ref() {
476            if let Err(e) = sender.send(ExecutionEvent::OrderCanceledBatch(batch)) {
477                log::warn!("Failed to send order canceled batch: {e}");
478            }
479        } else {
480            log::warn!("Cannot send order canceled batch: sender not initialized");
481        }
482    }
483
484    /// Emits an account state event.
485    pub fn send_account_state(&self, state: AccountState) {
486        if let Err(e) = self.try_send_account_state(state) {
487            log::warn!("{e}");
488        }
489    }
490
491    /// Emits an account state event and returns any channel error to the caller.
492    ///
493    /// # Errors
494    ///
495    /// Returns an error if the sender is uninitialized or its receiver is closed.
496    pub fn try_send_account_state(&self, state: AccountState) -> anyhow::Result<()> {
497        let sender = self.sender.load();
498        let sender = sender
499            .as_ref()
500            .ok_or_else(|| anyhow::anyhow!("Cannot send account state: sender not initialized"))?;
501        sender
502            .send(ExecutionEvent::Account(state))
503            .map_err(|e| anyhow::anyhow!("Failed to send account state: {e}"))
504    }
505
506    /// Emits an execution report.
507    pub fn send_execution_report(&self, report: ExecutionReport) {
508        if let Err(e) = self.try_send_execution_report(report) {
509            log::warn!("{e}");
510        }
511    }
512
513    /// Emits an execution report and returns any channel error to the caller.
514    ///
515    /// # Errors
516    ///
517    /// Returns an error if the sender is not initialized or the receiving channel is closed.
518    pub fn try_send_execution_report(&self, report: ExecutionReport) -> anyhow::Result<()> {
519        let sender = self.sender.load();
520
521        let sender = sender.as_ref().ok_or_else(|| {
522            anyhow::anyhow!("Cannot send execution report: sender not initialized")
523        })?;
524
525        sender
526            .send(ExecutionEvent::Report(report))
527            .map_err(|e| anyhow::anyhow!("Failed to send execution report: {e}"))
528    }
529
530    /// Emits an order status report.
531    pub fn send_order_status_report(&self, report: OrderStatusReport) {
532        self.send_execution_report(ExecutionReport::Order(Box::new(report)));
533    }
534
535    /// Emits a fill report.
536    pub fn send_fill_report(&self, report: FillReport) {
537        self.send_execution_report(ExecutionReport::Fill(Box::new(report)));
538    }
539
540    /// Emits an order status report bundled with the fills that produced it.
541    pub fn send_order_with_fills(&self, report: OrderStatusReport, fills: Vec<FillReport>) {
542        self.send_execution_report(ExecutionReport::OrderWithFills(Box::new(report), fills));
543    }
544
545    /// Emits a position status report.
546    pub fn send_position_report(&self, report: PositionStatusReport) {
547        self.send_execution_report(ExecutionReport::Position(Box::new(report)));
548    }
549}
550
551#[cfg(test)]
552mod tests {
553    use nautilus_core::time::get_atomic_clock_static;
554    use nautilus_model::events::order::spec::OrderSubmittedSpec;
555    use rstest::rstest;
556
557    use super::*;
558
559    fn create_emitter() -> ExecutionEventEmitter {
560        ExecutionEventEmitter::new(
561            get_atomic_clock_static(),
562            TraderId::from("TRADER-001"),
563            AccountId::from("SIM-001"),
564            AccountType::Cash,
565            None,
566        )
567    }
568
569    fn create_order_event() -> OrderEventAny {
570        OrderEventAny::Submitted(
571            OrderSubmittedSpec::builder()
572                .client_order_id(ClientOrderId::from("O-001"))
573                .build(),
574        )
575    }
576
577    #[rstest]
578    fn test_clone_before_set_sender_observes_sender() {
579        let mut emitter = create_emitter();
580        let cloned = emitter.clone();
581        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
582
583        emitter.set_sender(tx);
584
585        assert!(cloned.is_initialized());
586        cloned.send_order_event(create_order_event());
587        assert!(matches!(
588            rx.try_recv(),
589            Ok(ExecutionEvent::Order(OrderEventAny::Submitted(_)))
590        ));
591    }
592
593    #[rstest]
594    fn test_set_sender_replaces_existing_sender() {
595        let mut emitter = create_emitter();
596        let (tx_a, mut rx_a) = tokio::sync::mpsc::unbounded_channel();
597        let (tx_b, mut rx_b) = tokio::sync::mpsc::unbounded_channel();
598
599        emitter.set_sender(tx_a);
600        emitter.set_sender(tx_b);
601        emitter.send_order_event(create_order_event());
602
603        assert!(matches!(
604            rx_a.try_recv(),
605            Err(tokio::sync::mpsc::error::TryRecvError::Disconnected)
606        ));
607        assert!(matches!(
608            rx_b.try_recv(),
609            Ok(ExecutionEvent::Order(OrderEventAny::Submitted(_)))
610        ));
611    }
612
613    #[rstest]
614    fn test_never_initialized_sender_drops_or_errors() {
615        let emitter = create_emitter();
616
617        emitter.send_order_event(create_order_event());
618        let error = emitter
619            .try_send_order_event(create_order_event())
620            .unwrap_err();
621
622        assert_eq!(
623            error.to_string(),
624            "Cannot send order event: sender not initialized"
625        );
626    }
627}