Skip to main content

nautilus_trading/python/
algorithm.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//! Python bindings for execution algorithms.
17
18use std::{cell::UnsafeCell, collections::HashMap, fmt::Debug, rc::Rc};
19
20use jiff::Timestamp;
21use nautilus_common::{
22    actor::{DataActor, DataActorNative, data_actor::DataActorCore},
23    component::Component,
24    enums::ComponentState,
25    messages::system::{QueueStateChanged, SocketStateChanged},
26    python::{
27        cache::PyCache,
28        clock::PyClock,
29        logging::{PyLogger, format_exception},
30        wrappers::get_python_message_bus,
31    },
32    runner::SystemChannel,
33    signal::Signal,
34    timer::TimeEvent,
35};
36use nautilus_core::{
37    UnixNanos,
38    python::{to_pyruntime_err, to_pyvalue_err, upgrade_py_weakref},
39};
40use nautilus_model::{
41    data::{CustomData, DataType},
42    enums::{TimeInForce, TriggerType},
43    events::{
44        OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
45        OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled, OrderInitialized,
46        OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
47        OrderSubmitted, OrderTriggered, OrderUpdated, PositionChanged, PositionClosed,
48        PositionEvent, PositionOpened,
49    },
50    identifiers::{ActorId, ClientId, ExecAlgorithmId, PositionId, TraderId},
51    orders::{LimitOrder, MarketOrder, MarketToLimitOrder, Order, OrderAny, OrderList},
52    python::{events::order::order_event_to_pyobject, orders::pyobject_to_order_any},
53    types::{Price, Quantity},
54};
55use nautilus_portfolio::python::PyPortfolio;
56use pyo3::{
57    IntoPyObjectExt,
58    prelude::*,
59    types::{PyDict, PyList, PyWeakrefReference},
60};
61use ustr::Ustr;
62
63use crate::algorithm::{
64    EmulatedOrderSubmissionError, ExecutionAlgorithm, ExecutionAlgorithmConfig,
65    ExecutionAlgorithmCore, ExecutionAlgorithmNative, ImportableExecutionAlgorithmConfig,
66};
67
68/// Inner state of `PyExecutionAlgorithm`, shared by the Python and Rust registries.
69pub struct PyExecutionAlgorithmInner {
70    core: ExecutionAlgorithmCore,
71    py_self: Option<Py<PyWeakrefReference>>,
72    config: Option<Py<PyAny>>,
73    logger: PyLogger,
74}
75
76impl PyExecutionAlgorithmInner {
77    // The trader owns the wrapper for as long as the algorithm stays registered, so a collected
78    // wrapper propagates as an error rather than a skipped callback.
79    fn python_instance(&self) -> PyResult<Option<Py<PyAny>>> {
80        upgrade_py_weakref(self.py_self.as_ref(), &self.core.exec_algorithm_id)
81    }
82}
83
84impl Debug for PyExecutionAlgorithmInner {
85    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
86        f.debug_struct(stringify!(PyExecutionAlgorithmInner))
87            .field("core", &self.core)
88            .field(
89                "py_self",
90                &self.py_self.as_ref().map(|_| "<Py<PyWeakrefReference>>"),
91            )
92            .field("config", &self.config.as_ref().map(|_| "<Py<PyAny>>"))
93            .field("logger", &self.logger)
94            .finish()
95    }
96}
97
98/// Python-facing wrapper for execution algorithms.
99#[allow(non_camel_case_types)]
100#[pyo3::pyclass(
101    module = "nautilus_trader.trading",
102    name = "ExecutionAlgorithm",
103    unsendable,
104    subclass,
105    skip_from_py_object,
106    weakref
107)]
108#[pyo3_stub_gen::derive::gen_stub_pyclass(module = "nautilus_trader.trading")]
109#[derive(Clone)]
110pub struct PyExecutionAlgorithm {
111    inner: Rc<UnsafeCell<PyExecutionAlgorithmInner>>,
112}
113
114impl Debug for PyExecutionAlgorithm {
115    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
116        f.debug_struct(stringify!(PyExecutionAlgorithm))
117            .field("inner", &self.inner())
118            .finish()
119    }
120}
121
122impl PyExecutionAlgorithm {
123    #[inline]
124    #[allow(unsafe_code)]
125    pub(crate) fn inner(&self) -> &PyExecutionAlgorithmInner {
126        // SAFETY: `PyExecutionAlgorithm` is `unsendable` so access is single-threaded, and
127        // callers never hold a mutable and shared reference simultaneously.
128        unsafe { &*self.inner.get() }
129    }
130
131    #[inline]
132    #[allow(unsafe_code, clippy::mut_from_ref)]
133    pub(crate) fn inner_mut(&self) -> &mut PyExecutionAlgorithmInner {
134        // SAFETY: `PyExecutionAlgorithm` is `unsendable` so access is single-threaded, and
135        // callers never hold a mutable and shared reference simultaneously.
136        unsafe { &mut *self.inner.get() }
137    }
138}
139
140impl PyExecutionAlgorithm {
141    /// Creates a new `PyExecutionAlgorithm` instance.
142    #[must_use]
143    pub fn new(config: Option<ExecutionAlgorithmConfig>) -> Self {
144        let mut config = config.unwrap_or_default();
145        if config.exec_algorithm_id.is_none() {
146            config.exec_algorithm_id = Some(ExecAlgorithmId::new(stringify!(ExecutionAlgorithm)));
147        }
148
149        let core = ExecutionAlgorithmCore::new(config);
150        let logger = PyLogger::new(core.actor.actor_id.as_str());
151
152        let inner = PyExecutionAlgorithmInner {
153            core,
154            py_self: None,
155            config: None,
156            logger,
157        };
158
159        Self {
160            inner: Rc::new(UnsafeCell::new(inner)),
161        }
162    }
163
164    /// Sets the Python instance reference for method dispatch.
165    ///
166    /// Only a weak reference is stored, so the caller keeps ownership of `py_obj`. The trader
167    /// owns registered wrappers; an unregistered execution algorithm stays collectable.
168    ///
169    /// # Errors
170    ///
171    /// Returns an error if `py_obj` cannot be weakly referenced.
172    pub fn set_python_instance(&mut self, py_obj: &Bound<'_, PyAny>) -> PyResult<()> {
173        self.inner_mut().py_self = Some(PyWeakrefReference::new(py_obj)?.unbind());
174        Ok(())
175    }
176
177    /// Stores the original Python config object passed at construction.
178    pub fn set_config(&mut self, config: Option<Py<PyAny>>) {
179        self.inner_mut().config = config;
180    }
181
182    /// Updates the runtime execution algorithm ID before registration.
183    pub fn set_exec_algorithm_id(&mut self, exec_algorithm_id: ExecAlgorithmId) {
184        let actor_id = ActorId::new(exec_algorithm_id.inner());
185        let inner = self.inner_mut();
186
187        inner.core.config.exec_algorithm_id = Some(exec_algorithm_id);
188        inner.core.exec_algorithm_id = exec_algorithm_id;
189        inner.core.actor.actor_id = actor_id;
190        inner.core.actor.config.actor_id = Some(actor_id);
191        inner.logger = PyLogger::new(inner.core.actor.actor_id.as_str());
192    }
193
194    /// Updates the runtime `log_events` setting.
195    pub fn set_log_events(&mut self, log_events: bool) {
196        let inner = self.inner_mut();
197        inner.core.config.log_events = log_events;
198        inner.core.actor.config.log_events = log_events;
199    }
200
201    /// Updates the runtime `log_commands` setting.
202    pub fn set_log_commands(&mut self, log_commands: bool) {
203        let inner = self.inner_mut();
204        inner.core.config.log_commands = log_commands;
205        inner.core.actor.config.log_commands = log_commands;
206    }
207
208    /// Returns the execution algorithm ID.
209    #[must_use]
210    pub fn exec_algorithm_id(&self) -> ExecAlgorithmId {
211        self.inner().core.exec_algorithm_id
212    }
213
214    fn dispatch_no_args(&self, method_name: &str) -> PyResult<()> {
215        if let Some(py_self) = self.inner().python_instance()? {
216            Python::attach(|py| py_self.call_method0(py, method_name))?;
217        }
218        Ok(())
219    }
220
221    fn dispatch_time_event(&self, event: &TimeEvent) -> PyResult<()> {
222        if let Some(py_self) = self.inner().python_instance()? {
223            Python::attach(|py| {
224                py_self.call_method1(py, "on_time_event", (event.clone().into_py_any(py)?,))
225            })?;
226        }
227        Ok(())
228    }
229
230    fn dispatch_on_signal(&self, signal: &Signal) -> PyResult<()> {
231        if let Some(py_self) = self.inner().python_instance()? {
232            Python::attach(|py| {
233                py_self.call_method1(py, "on_signal", (signal.clone().into_py_any(py)?,))
234            })?;
235        }
236        Ok(())
237    }
238
239    fn dispatch_on_queue_state(&self, event: &QueueStateChanged) -> PyResult<()> {
240        if let Some(py_self) = self.inner().python_instance()? {
241            Python::attach(|py| {
242                py_self.call_method1(py, "on_queue_state", (event.clone().into_py_any(py)?,))
243            })?;
244        }
245        Ok(())
246    }
247
248    fn dispatch_on_socket_state(&self, event: &SocketStateChanged) -> PyResult<()> {
249        if let Some(py_self) = self.inner().python_instance()? {
250            Python::attach(|py| {
251                py_self.call_method1(py, "on_socket_state", (event.clone().into_py_any(py)?,))
252            })?;
253        }
254        Ok(())
255    }
256
257    fn dispatch_on_order(&self, order: OrderAny) -> PyResult<()> {
258        if let Some(py_self) = self.inner().python_instance()? {
259            Python::attach(|py| {
260                let py_order = nautilus_model::python::orders::order_any_to_pyobject(py, order)?;
261                py_self.call_method1(py, "on_order", (py_order,))
262            })?;
263        }
264        Ok(())
265    }
266
267    fn dispatch_on_order_list(&self, order_list: OrderList, orders: Vec<OrderAny>) -> PyResult<()> {
268        if let Some(py_self) = self.inner().python_instance()? {
269            Python::attach(|py| -> PyResult<()> {
270                let py_order_list = order_list.into_py_any(py)?;
271                let py_orders: Vec<_> = orders
272                    .into_iter()
273                    .map(|order| nautilus_model::python::orders::order_any_to_pyobject(py, order))
274                    .collect::<PyResult<Vec<_>>>()?;
275                let py_orders = PyList::new(py, py_orders)?;
276                py_self.call_method1(py, "on_order_list", (py_order_list, py_orders))?;
277                Ok(())
278            })?;
279        }
280        Ok(())
281    }
282
283    fn has_python_override(&self, method_name: &str) -> PyResult<bool> {
284        Python::attach(|py| -> PyResult<bool> {
285            let Some(py_self) = self.inner().python_instance()? else {
286                return Ok(false);
287            };
288
289            let instance_type = py_self.bind(py).get_type();
290            let instance_method = instance_type.getattr(method_name)?;
291            let base_method = py.get_type::<Self>().getattr(method_name)?;
292
293            Ok(!instance_method.is(&base_method))
294        })
295    }
296
297    fn dispatch_order_event(&self, method_name: &str, event: OrderEventAny) -> PyResult<()> {
298        if let Some(py_self) = self.inner().python_instance()? {
299            Python::attach(|py| {
300                let py_event = order_event_to_pyobject(py, event)?;
301                py_self.call_method1(py, method_name, (py_event,))
302            })?;
303        }
304        Ok(())
305    }
306
307    fn dispatch_position_event(&self, method_name: &str, event: PositionEvent) -> PyResult<()> {
308        if let Some(py_self) = self.inner().python_instance()? {
309            Python::attach(|py| {
310                let py_event = match event {
311                    PositionEvent::PositionOpened(event) => event.into_py_any(py)?,
312                    PositionEvent::PositionChanged(event) => event.into_py_any(py)?,
313                    PositionEvent::PositionClosed(event) => event.into_py_any(py)?,
314                    PositionEvent::PositionAdjusted(event) => event.into_py_any(py)?,
315                };
316                py_self.call_method1(py, method_name, (py_event,))
317            })?;
318        }
319        Ok(())
320    }
321
322    fn tags_to_ustr(tags: Option<Vec<String>>) -> Option<Vec<Ustr>> {
323        tags.map(|tags| tags.into_iter().map(|s| Ustr::from(&s)).collect())
324    }
325
326    fn primary_order_for_spawn(
327        &self,
328        py: Python<'_>,
329        primary: Py<PyAny>,
330        quantity: Quantity,
331        reduce_primary: bool,
332    ) -> PyResult<OrderAny> {
333        if !self.inner().core.actor.is_registered() {
334            return Err(to_pyruntime_err(
335                "ExecutionAlgorithm must be registered before spawning orders",
336            ));
337        }
338
339        let primary = pyobject_to_order_any(py, primary)?;
340        let cached_primary = {
341            let cache = self.inner().core.actor.cache_ref();
342            cache
343                .order(&primary.client_order_id())
344                .map(|order| order.clone())
345        };
346
347        let primary = if reduce_primary {
348            cached_primary.ok_or_else(|| {
349                to_pyruntime_err(format!(
350                    "Cannot reduce primary order {}: order not found in cache",
351                    primary.client_order_id()
352                ))
353            })?
354        } else {
355            cached_primary.unwrap_or(primary)
356        };
357
358        if reduce_primary && quantity > primary.leaves_qty() {
359            return Err(to_pyvalue_err(format!(
360                "Spawn quantity {quantity} exceeds primary leaves_qty {}",
361                primary.leaves_qty()
362            )));
363        }
364
365        if reduce_primary && primary.is_closed() {
366            return Err(to_pyvalue_err(format!(
367                "Cannot reduce closed primary order {}",
368                primary.client_order_id()
369            )));
370        }
371
372        Ok(primary)
373    }
374}
375
376impl DataActorNative for PyExecutionAlgorithm {
377    fn core(&self) -> &DataActorCore {
378        DataActorNative::core(&self.inner().core)
379    }
380
381    fn core_mut(&mut self) -> &mut DataActorCore {
382        DataActorNative::core_mut(&mut self.inner_mut().core)
383    }
384}
385
386impl ExecutionAlgorithmNative for PyExecutionAlgorithm {
387    fn exec_algorithm_core(&self) -> &ExecutionAlgorithmCore {
388        &self.inner().core
389    }
390
391    fn exec_algorithm_core_mut(&mut self) -> &mut ExecutionAlgorithmCore {
392        &mut self.inner_mut().core
393    }
394}
395
396impl ExecutionAlgorithm for PyExecutionAlgorithm {
397    fn on_order(&mut self, order: OrderAny) -> anyhow::Result<()> {
398        self.dispatch_on_order(order)
399            .map_err(|e| anyhow::anyhow!("Python on_order failed:\n{}", format_exception(&e)))
400    }
401
402    fn on_order_list(
403        &mut self,
404        order_list: OrderList,
405        orders: Vec<OrderAny>,
406    ) -> anyhow::Result<()> {
407        if self.has_python_override("on_order_list").map_err(|e| {
408            anyhow::anyhow!("Python override lookup failed:\n{}", format_exception(&e))
409        })? {
410            return self
411                .dispatch_on_order_list(order_list, orders)
412                .map_err(|e| {
413                    anyhow::anyhow!("Python on_order_list failed:\n{}", format_exception(&e))
414                });
415        }
416
417        for order in orders {
418            self.on_order(order)?;
419        }
420        Ok(())
421    }
422
423    fn on_start(&mut self) -> anyhow::Result<()> {
424        log::info!("Starting {}", self.exec_algorithm_id());
425        Ok(())
426    }
427
428    fn on_stop(&mut self) -> anyhow::Result<()> {
429        Ok(())
430    }
431
432    fn on_reset(&mut self) -> anyhow::Result<()> {
433        self.unsubscribe_all_strategy_events();
434        self.inner_mut().core.reset();
435        Ok(())
436    }
437
438    fn on_time_event(&mut self, _event: &TimeEvent) -> anyhow::Result<()> {
439        Ok(())
440    }
441
442    fn on_order_initialized(&mut self, event: OrderInitialized) {
443        let result =
444            self.dispatch_order_event("on_order_initialized", OrderEventAny::Initialized(event));
445        self.inner()
446            .logger
447            .log_callback_error("on_order_initialized", result);
448    }
449
450    fn on_order_denied(&mut self, event: OrderDenied) {
451        let result = self.dispatch_order_event("on_order_denied", OrderEventAny::Denied(event));
452        self.inner()
453            .logger
454            .log_callback_error("on_order_denied", result);
455    }
456
457    fn on_order_emulated(&mut self, event: OrderEmulated) {
458        let result = self.dispatch_order_event("on_order_emulated", OrderEventAny::Emulated(event));
459        self.inner()
460            .logger
461            .log_callback_error("on_order_emulated", result);
462    }
463
464    fn on_order_released(&mut self, event: OrderReleased) {
465        let result = self.dispatch_order_event("on_order_released", OrderEventAny::Released(event));
466        self.inner()
467            .logger
468            .log_callback_error("on_order_released", result);
469    }
470
471    fn on_order_submitted(&mut self, event: OrderSubmitted) {
472        let result =
473            self.dispatch_order_event("on_order_submitted", OrderEventAny::Submitted(event));
474        self.inner()
475            .logger
476            .log_callback_error("on_order_submitted", result);
477    }
478
479    fn on_order_rejected(&mut self, event: OrderRejected) {
480        let result = self.dispatch_order_event("on_order_rejected", OrderEventAny::Rejected(event));
481        self.inner()
482            .logger
483            .log_callback_error("on_order_rejected", result);
484    }
485
486    fn on_order_accepted(&mut self, event: OrderAccepted) {
487        let result = self.dispatch_order_event("on_order_accepted", OrderEventAny::Accepted(event));
488        self.inner()
489            .logger
490            .log_callback_error("on_order_accepted", result);
491    }
492
493    fn on_algo_order_canceled(&mut self, event: OrderCanceled) {
494        let result = self.dispatch_order_event("on_order_canceled", OrderEventAny::Canceled(event));
495        self.inner()
496            .logger
497            .log_callback_error("on_order_canceled", result);
498    }
499
500    fn on_order_expired(&mut self, event: OrderExpired) {
501        let result = self.dispatch_order_event("on_order_expired", OrderEventAny::Expired(event));
502        self.inner()
503            .logger
504            .log_callback_error("on_order_expired", result);
505    }
506
507    fn on_order_triggered(&mut self, event: OrderTriggered) {
508        let result =
509            self.dispatch_order_event("on_order_triggered", OrderEventAny::Triggered(event));
510        self.inner()
511            .logger
512            .log_callback_error("on_order_triggered", result);
513    }
514
515    fn on_order_pending_update(&mut self, event: OrderPendingUpdate) {
516        let result = self.dispatch_order_event(
517            "on_order_pending_update",
518            OrderEventAny::PendingUpdate(event),
519        );
520        self.inner()
521            .logger
522            .log_callback_error("on_order_pending_update", result);
523    }
524
525    fn on_order_pending_cancel(&mut self, event: OrderPendingCancel) {
526        let result = self.dispatch_order_event(
527            "on_order_pending_cancel",
528            OrderEventAny::PendingCancel(event),
529        );
530        self.inner()
531            .logger
532            .log_callback_error("on_order_pending_cancel", result);
533    }
534
535    fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {
536        let result = self.dispatch_order_event(
537            "on_order_modify_rejected",
538            OrderEventAny::ModifyRejected(event),
539        );
540        self.inner()
541            .logger
542            .log_callback_error("on_order_modify_rejected", result);
543    }
544
545    fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {
546        let result = self.dispatch_order_event(
547            "on_order_cancel_rejected",
548            OrderEventAny::CancelRejected(event),
549        );
550        self.inner()
551            .logger
552            .log_callback_error("on_order_cancel_rejected", result);
553    }
554
555    fn on_order_updated(&mut self, event: OrderUpdated) {
556        let result = self.dispatch_order_event("on_order_updated", OrderEventAny::Updated(event));
557        self.inner()
558            .logger
559            .log_callback_error("on_order_updated", result);
560    }
561
562    fn on_algo_order_filled(&mut self, event: OrderFilled) {
563        let result = self.dispatch_order_event("on_order_filled", OrderEventAny::Filled(event));
564        self.inner()
565            .logger
566            .log_callback_error("on_order_filled", result);
567    }
568
569    fn on_order_fill_voided(&mut self, event: &OrderFillVoided) {
570        let result = self.dispatch_order_event(
571            "on_order_fill_voided",
572            OrderEventAny::FillVoided(event.clone()),
573        );
574        self.inner()
575            .logger
576            .log_callback_error("on_order_fill_voided", result);
577    }
578
579    fn on_order_event(&mut self, event: OrderEventAny) {
580        let result = self.dispatch_order_event("on_order_event", event);
581        self.inner()
582            .logger
583            .log_callback_error("on_order_event", result);
584    }
585
586    fn on_position_opened(&mut self, event: PositionOpened) {
587        let result = self
588            .dispatch_position_event("on_position_opened", PositionEvent::PositionOpened(event));
589        self.inner()
590            .logger
591            .log_callback_error("on_position_opened", result);
592    }
593
594    fn on_position_changed(&mut self, event: PositionChanged) {
595        let result = self
596            .dispatch_position_event("on_position_changed", PositionEvent::PositionChanged(event));
597        self.inner()
598            .logger
599            .log_callback_error("on_position_changed", result);
600    }
601
602    fn on_position_closed(&mut self, event: PositionClosed) {
603        let result = self
604            .dispatch_position_event("on_position_closed", PositionEvent::PositionClosed(event));
605        self.inner()
606            .logger
607            .log_callback_error("on_position_closed", result);
608    }
609
610    fn on_position_event(&mut self, event: PositionEvent) {
611        let result = self.dispatch_position_event("on_position_event", event);
612        self.inner()
613            .logger
614            .log_callback_error("on_position_event", result);
615    }
616}
617
618impl DataActor for PyExecutionAlgorithm {
619    fn on_start(&mut self) -> anyhow::Result<()> {
620        ExecutionAlgorithm::on_start(self)?;
621        self.dispatch_no_args("on_start")
622            .map_err(|e| anyhow::anyhow!("Python on_start failed:\n{}", format_exception(&e)))
623    }
624
625    fn on_stop(&mut self) -> anyhow::Result<()> {
626        ExecutionAlgorithm::on_stop(self)?;
627        self.dispatch_no_args("on_stop")
628            .map_err(|e| anyhow::anyhow!("Python on_stop failed:\n{}", format_exception(&e)))
629    }
630
631    fn on_resume(&mut self) -> anyhow::Result<()> {
632        ExecutionAlgorithm::on_resume(self)?;
633        self.dispatch_no_args("on_resume")
634            .map_err(|e| anyhow::anyhow!("Python on_resume failed:\n{}", format_exception(&e)))
635    }
636
637    fn on_reset(&mut self) -> anyhow::Result<()> {
638        ExecutionAlgorithm::on_reset(self)?;
639        self.dispatch_no_args("on_reset")
640            .map_err(|e| anyhow::anyhow!("Python on_reset failed:\n{}", format_exception(&e)))
641    }
642
643    fn on_dispose(&mut self) -> anyhow::Result<()> {
644        self.dispatch_no_args("on_dispose")
645            .map_err(|e| anyhow::anyhow!("Python on_dispose failed:\n{}", format_exception(&e)))
646    }
647
648    fn on_degrade(&mut self) -> anyhow::Result<()> {
649        self.dispatch_no_args("on_degrade")
650            .map_err(|e| anyhow::anyhow!("Python on_degrade failed:\n{}", format_exception(&e)))
651    }
652
653    fn on_fault(&mut self) -> anyhow::Result<()> {
654        self.dispatch_no_args("on_fault")
655            .map_err(|e| anyhow::anyhow!("Python on_fault failed:\n{}", format_exception(&e)))
656    }
657
658    fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()> {
659        ExecutionAlgorithm::on_time_event(self, event)?;
660        self.dispatch_time_event(event)
661            .map_err(|e| anyhow::anyhow!("Python on_time_event failed:\n{}", format_exception(&e)))
662    }
663
664    fn on_signal(&mut self, signal: &Signal) -> anyhow::Result<()> {
665        self.dispatch_on_signal(signal)
666            .map_err(|e| anyhow::anyhow!("Python on_signal failed:\n{}", format_exception(&e)))
667    }
668
669    fn on_queue_state(&mut self, event: &QueueStateChanged) -> anyhow::Result<()> {
670        self.dispatch_on_queue_state(event)
671            .map_err(|e| anyhow::anyhow!("Python on_queue_state failed:\n{}", format_exception(&e)))
672    }
673
674    fn on_socket_state(&mut self, event: &SocketStateChanged) -> anyhow::Result<()> {
675        self.dispatch_on_socket_state(event).map_err(|e| {
676            anyhow::anyhow!("Python on_socket_state failed:\n{}", format_exception(&e))
677        })
678    }
679}
680
681#[pyo3::pymethods]
682#[pyo3_stub_gen::derive::gen_stub_pymethods]
683#[allow(
684    clippy::large_types_passed_by_value,
685    reason = "PyO3 callbacks accept Python-owned event values"
686)]
687#[expect(
688    clippy::unused_self,
689    reason = "default PyO3 callbacks must remain instance methods"
690)]
691impl PyExecutionAlgorithm {
692    /// Creates a new [`PyExecutionAlgorithm`] instance.
693    #[new]
694    #[pyo3(signature = (config=None))]
695    fn py_new(config: Option<Py<PyAny>>) -> Self {
696        let algorithm_config = config
697            .as_ref()
698            .and_then(|obj| Python::attach(|py| obj.extract::<ExecutionAlgorithmConfig>(py).ok()));
699        let mut algorithm = Self::new(algorithm_config);
700        algorithm.set_config(config);
701        algorithm
702    }
703
704    /// Captures the Python self reference for Rust→Python event dispatch.
705    #[pyo3(signature = (config=None))]
706    fn __init__(slf: &Bound<'_, Self>, config: Option<Py<PyAny>>) -> PyResult<()> {
707        let retained_config = if config.is_none() {
708            Python::attach(|py| {
709                slf.borrow()
710                    .inner()
711                    .config
712                    .as_ref()
713                    .map(|config| config.clone_ref(py))
714            })
715        } else {
716            None
717        };
718        let has_configured_id = if let Some(config) = config.as_ref().or(retained_config.as_ref()) {
719            Python::attach(|py| slf.borrow_mut().configure_from_py_config(config.bind(py)))
720                .map_err(to_pyvalue_err)?
721        } else {
722            false
723        };
724
725        if !has_configured_id {
726            let py_type = slf.get_type();
727            let type_name = py_type.name()?;
728            let exec_algorithm_id =
729                ExecAlgorithmId::new_checked(type_name.to_str()?).map_err(to_pyvalue_err)?;
730            slf.borrow_mut().set_exec_algorithm_id(exec_algorithm_id);
731        }
732        let mut borrowed = slf.borrow_mut();
733        borrowed.set_python_instance(slf.as_any())?;
734        if config.is_some() {
735            borrowed.set_config(config);
736        }
737        Ok(())
738    }
739
740    #[getter]
741    #[pyo3(name = "trader_id")]
742    fn py_trader_id(&self) -> Option<TraderId> {
743        self.inner().core.actor.trader_id()
744    }
745
746    #[getter]
747    #[pyo3(name = "exec_algorithm_id")]
748    fn py_exec_algorithm_id(&self) -> ExecAlgorithmId {
749        self.exec_algorithm_id()
750    }
751
752    #[getter]
753    #[pyo3(name = "config")]
754    fn py_config(&self, py: Python<'_>) -> Option<Py<PyAny>> {
755        self.inner()
756            .config
757            .as_ref()
758            .map(|config| config.clone_ref(py))
759    }
760
761    /// Returns an importable configuration for this execution algorithm.
762    #[pyo3(name = "to_importable_config")]
763    fn py_to_importable_config(
764        &self,
765        py: Python<'_>,
766    ) -> PyResult<ImportableExecutionAlgorithmConfig> {
767        let py_self = self
768            .inner()
769            .python_instance()?
770            .ok_or_else(|| to_pyruntime_err("Python execution algorithm instance is not set"))?;
771        let exec_algorithm_path = py_type_path(py_self.bind(py))?;
772
773        let Some(config) = self.inner().config.as_ref() else {
774            return Ok(ImportableExecutionAlgorithmConfig {
775                exec_algorithm_path,
776                config_path: String::new(),
777                config: HashMap::new(),
778            });
779        };
780        let config = config.bind(py);
781
782        Ok(ImportableExecutionAlgorithmConfig {
783            exec_algorithm_path,
784            config_path: py_type_path(config)?,
785            config: py_config_to_json(config)?,
786        })
787    }
788
789    #[getter]
790    #[pyo3(name = "clock")]
791    fn py_clock(&self) -> Option<PyClock> {
792        self.inner()
793            .core
794            .actor
795            .is_registered()
796            .then(|| PyClock::from_rc(self.inner().core.actor.clock_rc()))
797    }
798
799    #[getter]
800    #[pyo3(name = "cache")]
801    fn py_cache(&self) -> Option<PyCache> {
802        self.inner()
803            .core
804            .actor
805            .is_registered()
806            .then(|| PyCache::from_rc(self.inner().core.actor.cache_rc()))
807    }
808
809    #[getter]
810    #[pyo3(name = "portfolio")]
811    fn py_portfolio(&self) -> PyResult<PyPortfolio> {
812        if self.inner().core.actor.is_registered() {
813            Ok(PyPortfolio::from_rc(self.portfolio_rc()))
814        } else {
815            Err(to_pyruntime_err(
816                "ExecutionAlgorithm must be registered with a trader before accessing portfolio",
817            ))
818        }
819    }
820
821    #[getter]
822    #[pyo3(name = "log")]
823    fn py_log(&self) -> PyLogger {
824        self.inner().logger.clone()
825    }
826
827    #[getter]
828    #[pyo3(name = "state")]
829    fn py_state(&self) -> ComponentState {
830        self.inner().core.actor.state()
831    }
832
833    #[pyo3(name = "is_registered")]
834    fn py_is_registered(&self) -> bool {
835        self.inner().core.actor.is_registered()
836    }
837
838    #[pyo3(name = "is_ready")]
839    fn py_is_ready(&self) -> bool {
840        Component::is_ready(self)
841    }
842
843    #[pyo3(name = "is_running")]
844    fn py_is_running(&self) -> bool {
845        Component::is_running(self)
846    }
847
848    #[pyo3(name = "is_stopped")]
849    fn py_is_stopped(&self) -> bool {
850        Component::is_stopped(self)
851    }
852
853    #[pyo3(name = "is_disposed")]
854    fn py_is_disposed(&self) -> bool {
855        Component::is_disposed(self)
856    }
857
858    #[pyo3(name = "is_degraded")]
859    fn py_is_degraded(&self) -> bool {
860        Component::is_degraded(self)
861    }
862
863    #[pyo3(name = "is_faulted")]
864    fn py_is_faulted(&self) -> bool {
865        Component::is_faulted(self)
866    }
867
868    #[pyo3(name = "start")]
869    fn py_start(slf: PyRef<'_, Self>) -> PyResult<()> {
870        let mut exec_algorithm = slf.clone();
871        drop(slf);
872        Component::start(&mut exec_algorithm).map_err(to_pyruntime_err)
873    }
874
875    #[pyo3(name = "stop")]
876    fn py_stop(slf: PyRef<'_, Self>) -> PyResult<()> {
877        let mut exec_algorithm = slf.clone();
878        drop(slf);
879        Component::stop(&mut exec_algorithm).map_err(to_pyruntime_err)
880    }
881
882    #[pyo3(name = "resume")]
883    fn py_resume(slf: PyRef<'_, Self>) -> PyResult<()> {
884        let mut exec_algorithm = slf.clone();
885        drop(slf);
886        Component::resume(&mut exec_algorithm).map_err(to_pyruntime_err)
887    }
888
889    #[pyo3(name = "reset")]
890    fn py_reset(slf: PyRef<'_, Self>) -> PyResult<()> {
891        let mut exec_algorithm = slf.clone();
892        drop(slf);
893        Component::reset(&mut exec_algorithm).map_err(to_pyruntime_err)
894    }
895
896    #[pyo3(name = "dispose")]
897    fn py_dispose(slf: PyRef<'_, Self>) -> PyResult<()> {
898        let mut exec_algorithm = slf.clone();
899        drop(slf);
900        Component::dispose(&mut exec_algorithm).map_err(to_pyruntime_err)
901    }
902
903    #[pyo3(name = "degrade")]
904    fn py_degrade(slf: PyRef<'_, Self>) -> PyResult<()> {
905        let mut exec_algorithm = slf.clone();
906        drop(slf);
907        Component::degrade(&mut exec_algorithm).map_err(to_pyruntime_err)
908    }
909
910    #[pyo3(name = "fault")]
911    fn py_fault(slf: PyRef<'_, Self>) -> PyResult<()> {
912        let mut exec_algorithm = slf.clone();
913        drop(slf);
914        Component::fault(&mut exec_algorithm).map_err(to_pyruntime_err)
915    }
916
917    #[pyo3(name = "publish_data")]
918    fn py_publish_data(&self, data_type: &DataType, data: &CustomData) -> PyResult<()> {
919        self.ensure_registered_for_data()?;
920        DataActor::publish_data(self, data_type, data);
921        Ok(())
922    }
923
924    #[pyo3(name = "publish_signal")]
925    #[pyo3(signature = (name, value, ts_event=0))]
926    #[expect(
927        clippy::needless_pass_by_value,
928        reason = "PyO3 accepts an owned PyAny handle for Python signal values"
929    )]
930    fn py_publish_signal(
931        &self,
932        py: Python<'_>,
933        name: &str,
934        value: Py<PyAny>,
935        ts_event: u64,
936    ) -> PyResult<()> {
937        self.ensure_registered_for_data()?;
938        let value_str: String = value.bind(py).str()?.extract()?;
939        DataActor::publish_signal(self, name, value_str, UnixNanos::from(ts_event));
940        Ok(())
941    }
942
943    #[pyo3(name = "subscribe_signal")]
944    #[pyo3(signature = (name="", priority=None))]
945    fn py_subscribe_signal(&mut self, name: &str, priority: Option<u32>) -> PyResult<()> {
946        self.ensure_registered()?;
947        DataActor::subscribe_signal(self, name, priority);
948        Ok(())
949    }
950
951    #[pyo3(name = "subscribe_queue_state")]
952    #[pyo3(signature = (channel=None, priority=None))]
953    fn py_subscribe_queue_state(
954        &mut self,
955        channel: Option<SystemChannel>,
956        priority: Option<u32>,
957    ) -> PyResult<()> {
958        self.ensure_registered()?;
959        DataActor::subscribe_queue_state(self, channel, priority);
960        Ok(())
961    }
962
963    #[pyo3(name = "subscribe_socket_state")]
964    #[pyo3(signature = (client_id=None, endpoint=None, priority=None))]
965    fn py_subscribe_socket_state(
966        &mut self,
967        client_id: Option<ClientId>,
968        endpoint: Option<&str>,
969        priority: Option<u32>,
970    ) -> PyResult<()> {
971        self.ensure_registered()?;
972        DataActor::subscribe_socket_state(self, client_id, endpoint, priority);
973        Ok(())
974    }
975
976    #[pyo3(name = "unsubscribe_signal")]
977    #[pyo3(signature = (name=""))]
978    fn py_unsubscribe_signal(&mut self, name: &str) -> PyResult<()> {
979        self.ensure_registered()?;
980        DataActor::unsubscribe_signal(self, name);
981        Ok(())
982    }
983
984    #[pyo3(name = "unsubscribe_queue_state")]
985    #[pyo3(signature = (channel=None))]
986    fn py_unsubscribe_queue_state(&mut self, channel: Option<SystemChannel>) -> PyResult<()> {
987        self.ensure_registered()?;
988        DataActor::unsubscribe_queue_state(self, channel);
989        Ok(())
990    }
991
992    #[pyo3(name = "unsubscribe_socket_state")]
993    #[pyo3(signature = (client_id=None, endpoint=None))]
994    fn py_unsubscribe_socket_state(
995        &mut self,
996        client_id: Option<ClientId>,
997        endpoint: Option<&str>,
998    ) -> PyResult<()> {
999        self.ensure_registered()?;
1000        DataActor::unsubscribe_socket_state(self, client_id, endpoint);
1001        Ok(())
1002    }
1003
1004    #[pyo3(name = "on_start")]
1005    fn py_on_start(&mut self) {}
1006
1007    #[pyo3(name = "on_stop")]
1008    fn py_on_stop(&mut self) {}
1009
1010    #[pyo3(name = "on_resume")]
1011    fn py_on_resume(&mut self) {}
1012
1013    #[pyo3(name = "on_reset")]
1014    fn py_on_reset(&mut self) {}
1015
1016    #[pyo3(name = "on_dispose")]
1017    fn py_on_dispose(&mut self) {}
1018
1019    #[pyo3(name = "on_degrade")]
1020    fn py_on_degrade(&mut self) {}
1021
1022    #[pyo3(name = "on_fault")]
1023    fn py_on_fault(&mut self) {}
1024
1025    #[allow(unused_variables, clippy::needless_pass_by_value)]
1026    #[pyo3(name = "on_time_event")]
1027    fn py_on_time_event(&mut self, event: TimeEvent) {}
1028
1029    #[allow(unused_variables)]
1030    #[pyo3(name = "on_signal")]
1031    fn py_on_signal(&mut self, signal: &Signal) {}
1032
1033    #[allow(unused_variables, clippy::needless_pass_by_value)]
1034    #[pyo3(name = "on_queue_state")]
1035    fn py_on_queue_state(&mut self, event: QueueStateChanged) {}
1036
1037    #[allow(unused_variables, clippy::needless_pass_by_value)]
1038    #[pyo3(name = "on_socket_state")]
1039    fn py_on_socket_state(&mut self, event: SocketStateChanged) {}
1040
1041    #[allow(clippy::needless_pass_by_value)]
1042    #[pyo3(name = "execute")]
1043    fn py_execute(&mut self, command: Py<PyAny>) -> PyResult<()> {
1044        let _ = command;
1045        Err(to_pyruntime_err(
1046            "ExecutionAlgorithm.execute is invoked by the v2 runtime endpoint",
1047        ))
1048    }
1049
1050    #[allow(unused_variables, clippy::needless_pass_by_value)]
1051    #[pyo3(name = "on_order")]
1052    fn py_on_order(&mut self, order: Py<PyAny>) {}
1053
1054    #[allow(unused_variables, clippy::needless_pass_by_value)]
1055    #[pyo3(name = "on_order_list")]
1056    fn py_on_order_list(&mut self, order_list: Py<PyAny>, orders: Py<PyAny>) {}
1057
1058    #[pyo3(name = "spawn_market")]
1059    #[pyo3(signature = (
1060        primary,
1061        quantity,
1062        time_in_force = TimeInForce::Gtc,
1063        reduce_only = false,
1064        tags = None,
1065        reduce_primary = true
1066    ))]
1067    #[expect(clippy::too_many_arguments)]
1068    fn py_spawn_market(
1069        &mut self,
1070        py: Python<'_>,
1071        primary: Py<PyAny>,
1072        quantity: Quantity,
1073        time_in_force: TimeInForce,
1074        reduce_only: bool,
1075        tags: Option<Vec<String>>,
1076        reduce_primary: bool,
1077    ) -> PyResult<MarketOrder> {
1078        let mut primary = self.primary_order_for_spawn(py, primary, quantity, reduce_primary)?;
1079        Ok(ExecutionAlgorithm::spawn_market(
1080            self,
1081            &mut primary,
1082            quantity,
1083            time_in_force,
1084            reduce_only,
1085            Self::tags_to_ustr(tags),
1086            reduce_primary,
1087        ))
1088    }
1089
1090    #[pyo3(name = "spawn_limit")]
1091    #[pyo3(signature = (
1092        primary,
1093        quantity,
1094        price,
1095        time_in_force = TimeInForce::Gtc,
1096        expire_time = None,
1097        post_only = false,
1098        reduce_only = false,
1099        display_qty = None,
1100        emulation_trigger = None,
1101        tags = None,
1102        reduce_primary = true
1103    ))]
1104    #[expect(clippy::too_many_arguments)]
1105    fn py_spawn_limit(
1106        &mut self,
1107        py: Python<'_>,
1108        primary: Py<PyAny>,
1109        quantity: Quantity,
1110        price: Price,
1111        time_in_force: TimeInForce,
1112        expire_time: Option<Timestamp>,
1113        post_only: bool,
1114        reduce_only: bool,
1115        display_qty: Option<Quantity>,
1116        emulation_trigger: Option<TriggerType>,
1117        tags: Option<Vec<String>>,
1118        reduce_primary: bool,
1119    ) -> PyResult<LimitOrder> {
1120        let mut primary = self.primary_order_for_spawn(py, primary, quantity, reduce_primary)?;
1121        Ok(ExecutionAlgorithm::spawn_limit(
1122            self,
1123            &mut primary,
1124            quantity,
1125            price,
1126            time_in_force,
1127            expire_time.map(UnixNanos::from),
1128            post_only,
1129            reduce_only,
1130            display_qty,
1131            emulation_trigger,
1132            Self::tags_to_ustr(tags),
1133            reduce_primary,
1134        ))
1135    }
1136
1137    #[pyo3(name = "spawn_market_to_limit")]
1138    #[pyo3(signature = (
1139        primary,
1140        quantity,
1141        time_in_force = TimeInForce::Gtc,
1142        expire_time = None,
1143        reduce_only = false,
1144        display_qty = None,
1145        emulation_trigger = None,
1146        tags = None,
1147        reduce_primary = true
1148    ))]
1149    #[expect(clippy::too_many_arguments)]
1150    fn py_spawn_market_to_limit(
1151        &mut self,
1152        py: Python<'_>,
1153        primary: Py<PyAny>,
1154        quantity: Quantity,
1155        time_in_force: TimeInForce,
1156        expire_time: Option<Timestamp>,
1157        reduce_only: bool,
1158        display_qty: Option<Quantity>,
1159        emulation_trigger: Option<TriggerType>,
1160        tags: Option<Vec<String>>,
1161        reduce_primary: bool,
1162    ) -> PyResult<MarketToLimitOrder> {
1163        let mut primary = self.primary_order_for_spawn(py, primary, quantity, reduce_primary)?;
1164        Ok(ExecutionAlgorithm::spawn_market_to_limit(
1165            self,
1166            &mut primary,
1167            quantity,
1168            time_in_force,
1169            expire_time.map(UnixNanos::from),
1170            reduce_only,
1171            display_qty,
1172            emulation_trigger,
1173            Self::tags_to_ustr(tags),
1174            reduce_primary,
1175        ))
1176    }
1177
1178    #[pyo3(name = "deny_order")]
1179    fn py_deny_order(&mut self, py: Python<'_>, order: Py<PyAny>, reason: &str) -> PyResult<()> {
1180        let order = pyobject_to_order_any(py, order)?;
1181        ExecutionAlgorithm::deny_order(self, &order, Ustr::from(reason)).map_err(to_pyruntime_err)
1182    }
1183
1184    #[pyo3(name = "submit_order")]
1185    #[pyo3(signature = (order, position_id=None, client_id=None))]
1186    fn py_submit_order(
1187        &mut self,
1188        py: Python<'_>,
1189        order: Py<PyAny>,
1190        position_id: Option<PositionId>,
1191        client_id: Option<ClientId>,
1192    ) -> PyResult<()> {
1193        let order = pyobject_to_order_any(py, order)?;
1194        ExecutionAlgorithm::submit_order(self, order, position_id, client_id).map_err(|e| {
1195            if e.downcast_ref::<EmulatedOrderSubmissionError>().is_some() {
1196                to_pyvalue_err(e)
1197            } else {
1198                to_pyruntime_err(e)
1199            }
1200        })
1201    }
1202
1203    #[pyo3(name = "modify_order")]
1204    #[pyo3(signature = (order, quantity=None, price=None, trigger_price=None, client_id=None))]
1205    fn py_modify_order(
1206        &mut self,
1207        py: Python<'_>,
1208        order: Py<PyAny>,
1209        quantity: Option<Quantity>,
1210        price: Option<Price>,
1211        trigger_price: Option<Price>,
1212        client_id: Option<ClientId>,
1213    ) -> PyResult<()> {
1214        let mut order = pyobject_to_order_any(py, order)?;
1215        ExecutionAlgorithm::modify_order(
1216            self,
1217            &mut order,
1218            quantity,
1219            price,
1220            trigger_price,
1221            client_id,
1222        )
1223        .map_err(to_pyruntime_err)
1224    }
1225
1226    #[pyo3(name = "modify_order_in_place")]
1227    #[pyo3(signature = (order, quantity=None, price=None, trigger_price=None))]
1228    fn py_modify_order_in_place(
1229        &mut self,
1230        py: Python<'_>,
1231        order: Py<PyAny>,
1232        quantity: Option<Quantity>,
1233        price: Option<Price>,
1234        trigger_price: Option<Price>,
1235    ) -> PyResult<()> {
1236        let mut order = pyobject_to_order_any(py, order)?;
1237        ExecutionAlgorithm::modify_order_in_place(self, &mut order, quantity, price, trigger_price)
1238            .map_err(to_pyruntime_err)
1239    }
1240
1241    #[pyo3(name = "cancel_order")]
1242    #[pyo3(signature = (order, client_id=None))]
1243    fn py_cancel_order(
1244        &mut self,
1245        py: Python<'_>,
1246        order: Py<PyAny>,
1247        client_id: Option<ClientId>,
1248    ) -> PyResult<()> {
1249        let mut order = pyobject_to_order_any(py, order)?;
1250        ExecutionAlgorithm::cancel_order(self, &mut order, client_id).map_err(to_pyruntime_err)
1251    }
1252
1253    #[allow(unused_variables, clippy::needless_pass_by_value)]
1254    #[pyo3(name = "on_order_initialized")]
1255    fn py_on_order_initialized(&mut self, event: OrderInitialized) {}
1256
1257    #[allow(unused_variables, clippy::needless_pass_by_value)]
1258    #[pyo3(name = "on_order_event")]
1259    fn py_on_order_event(&mut self, event: Py<PyAny>) {}
1260
1261    #[allow(unused_variables)]
1262    #[pyo3(name = "on_order_denied")]
1263    fn py_on_order_denied(&mut self, event: OrderDenied) {}
1264
1265    #[allow(unused_variables)]
1266    #[pyo3(name = "on_order_emulated")]
1267    fn py_on_order_emulated(&mut self, event: OrderEmulated) {}
1268
1269    #[allow(unused_variables)]
1270    #[pyo3(name = "on_order_released")]
1271    fn py_on_order_released(&mut self, event: OrderReleased) {}
1272
1273    #[allow(unused_variables)]
1274    #[pyo3(name = "on_order_submitted")]
1275    fn py_on_order_submitted(&mut self, event: OrderSubmitted) {}
1276
1277    #[allow(unused_variables)]
1278    #[pyo3(name = "on_order_rejected")]
1279    fn py_on_order_rejected(&mut self, event: OrderRejected) {}
1280
1281    #[allow(unused_variables)]
1282    #[pyo3(name = "on_order_accepted")]
1283    fn py_on_order_accepted(&mut self, event: OrderAccepted) {}
1284
1285    #[allow(unused_variables)]
1286    #[pyo3(name = "on_order_canceled")]
1287    fn py_on_order_canceled(&mut self, event: OrderCanceled) {}
1288
1289    #[allow(unused_variables)]
1290    #[pyo3(name = "on_order_expired")]
1291    fn py_on_order_expired(&mut self, event: OrderExpired) {}
1292
1293    #[allow(unused_variables)]
1294    #[pyo3(name = "on_order_triggered")]
1295    fn py_on_order_triggered(&mut self, event: OrderTriggered) {}
1296
1297    #[allow(unused_variables)]
1298    #[pyo3(name = "on_order_pending_update")]
1299    fn py_on_order_pending_update(&mut self, event: OrderPendingUpdate) {}
1300
1301    #[allow(unused_variables)]
1302    #[pyo3(name = "on_order_pending_cancel")]
1303    fn py_on_order_pending_cancel(&mut self, event: OrderPendingCancel) {}
1304
1305    #[allow(unused_variables)]
1306    #[pyo3(name = "on_order_modify_rejected")]
1307    fn py_on_order_modify_rejected(&mut self, event: OrderModifyRejected) {}
1308
1309    #[allow(unused_variables)]
1310    #[pyo3(name = "on_order_cancel_rejected")]
1311    fn py_on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {}
1312
1313    #[allow(unused_variables)]
1314    #[pyo3(name = "on_order_updated")]
1315    fn py_on_order_updated(&mut self, event: OrderUpdated) {}
1316
1317    #[allow(unused_variables, clippy::needless_pass_by_value)]
1318    #[pyo3(name = "on_order_filled")]
1319    fn py_on_order_filled(&mut self, event: OrderFilled) {}
1320
1321    #[allow(unused_variables, clippy::needless_pass_by_value)]
1322    #[pyo3(name = "on_order_fill_voided")]
1323    fn py_on_order_fill_voided(&mut self, event: OrderFillVoided) {}
1324
1325    #[allow(unused_variables, clippy::needless_pass_by_value)]
1326    #[pyo3(name = "on_position_opened")]
1327    fn py_on_position_opened(&mut self, event: PositionOpened) {}
1328
1329    #[allow(unused_variables, clippy::needless_pass_by_value)]
1330    #[pyo3(name = "on_position_event")]
1331    fn py_on_position_event(&mut self, event: Py<PyAny>) {}
1332
1333    #[allow(unused_variables, clippy::needless_pass_by_value)]
1334    #[pyo3(name = "on_position_changed")]
1335    fn py_on_position_changed(&mut self, event: PositionChanged) {}
1336
1337    #[allow(unused_variables, clippy::needless_pass_by_value)]
1338    #[pyo3(name = "on_position_closed")]
1339    fn py_on_position_closed(&mut self, event: PositionClosed) {}
1340}
1341
1342#[pyo3_stub_gen::derive::gen_stub_pymethods]
1343#[pyo3::pymethods]
1344impl PyExecutionAlgorithm {
1345    #[pyo3(name = "publish_message", signature = (topic, message))]
1346    fn py_publish_message(
1347        slf: &Bound<'_, Self>,
1348        topic: &str,
1349        #[gen_stub(override_type(type_repr = "object"))] message: Py<PyAny>,
1350    ) -> PyResult<()> {
1351        let messages = get_python_message_bus(slf.as_any())?;
1352        messages.publish_message(topic, message)
1353    }
1354
1355    #[pyo3(name = "subscribe_topic")]
1356    #[pyo3(signature = (topic, handler, priority=0))]
1357    fn py_subscribe_topic(
1358        slf: &Bound<'_, Self>,
1359        topic: &str,
1360        #[gen_stub(override_type(type_repr = "collections.abc.Callable[[object], None]", imports = ("collections.abc",)))]
1361        handler: Py<PyAny>,
1362        priority: u32,
1363    ) -> PyResult<()> {
1364        let messages = get_python_message_bus(slf.as_any())?;
1365        messages.subscribe_topic(slf.py(), topic, handler, priority)
1366    }
1367
1368    #[pyo3(name = "unsubscribe_topic", signature = (topic, handler))]
1369    fn py_unsubscribe_topic(
1370        slf: &Bound<'_, Self>,
1371        topic: &str,
1372        #[gen_stub(override_type(type_repr = "collections.abc.Callable[[object], None]", imports = ("collections.abc",)))]
1373        handler: &Bound<'_, PyAny>,
1374    ) -> PyResult<()> {
1375        let messages = get_python_message_bus(slf.as_any())?;
1376        messages.unsubscribe_topic(topic, handler)
1377    }
1378}
1379
1380impl PyExecutionAlgorithm {
1381    /// Applies Python configuration overrides.
1382    ///
1383    /// Returns whether the config supplied an execution algorithm ID.
1384    ///
1385    /// # Errors
1386    ///
1387    /// Returns an error if an ID has an unsupported type or invalid value.
1388    pub fn configure_from_py_config(&mut self, config: &Bound<'_, PyAny>) -> anyhow::Result<bool> {
1389        let id = config
1390            .getattr("exec_algorithm_id")
1391            .ok()
1392            .filter(|id| !id.is_none())
1393            .or_else(|| config.getattr("actor_id").ok().filter(|id| !id.is_none()));
1394        let has_id = if let Some(id) = id {
1395            let exec_algorithm_id = if let Ok(exec_algorithm_id) = id.extract::<ExecAlgorithmId>() {
1396                exec_algorithm_id
1397            } else if let Ok(actor_id) = id.extract::<ActorId>() {
1398                ExecAlgorithmId::new_checked(actor_id.inner())?
1399            } else if let Ok(id) = id.extract::<String>() {
1400                ExecAlgorithmId::new_checked(&id)?
1401            } else {
1402                anyhow::bail!("Invalid `exec_algorithm_id`/`actor_id` type");
1403            };
1404            self.set_exec_algorithm_id(exec_algorithm_id);
1405            true
1406        } else {
1407            false
1408        };
1409
1410        if let Ok(log_events) = config.getattr("log_events")
1411            && let Ok(log_events) = log_events.extract::<bool>()
1412        {
1413            self.set_log_events(log_events);
1414        }
1415
1416        if let Ok(log_commands) = config.getattr("log_commands")
1417            && let Ok(log_commands) = log_commands.extract::<bool>()
1418        {
1419            self.set_log_commands(log_commands);
1420        }
1421
1422        Ok(has_id)
1423    }
1424
1425    fn ensure_registered_for_data(&self) -> PyResult<()> {
1426        if self.inner().core.actor.is_registered() {
1427            Ok(())
1428        } else {
1429            Err(to_pyruntime_err(
1430                "ExecutionAlgorithm must be registered before publishing data",
1431            ))
1432        }
1433    }
1434
1435    fn ensure_registered(&self) -> PyResult<()> {
1436        if self.inner().core.actor.is_registered() {
1437            Ok(())
1438        } else {
1439            Err(to_pyruntime_err(
1440                "ExecutionAlgorithm must be registered before managing subscriptions",
1441            ))
1442        }
1443    }
1444}
1445
1446#[pyo3::pymethods]
1447#[pyo3_stub_gen::derive::gen_stub_pymethods]
1448impl ExecutionAlgorithmConfig {
1449    /// Configuration for an execution algorithm.
1450    #[new]
1451    #[pyo3(signature = (
1452        exec_algorithm_id=None,
1453        log_events=true,
1454        log_commands=true,
1455        **_kwargs
1456    ))]
1457    fn py_new(
1458        #[gen_stub(override_type(type_repr = "model.ExecAlgorithmId | str | None"))]
1459        exec_algorithm_id: Option<&Bound<'_, PyAny>>,
1460        log_events: bool,
1461        log_commands: bool,
1462        _kwargs: Option<&Bound<'_, PyDict>>,
1463    ) -> PyResult<Self> {
1464        let exec_algorithm_id = exec_algorithm_id
1465            .map(|value| -> PyResult<ExecAlgorithmId> {
1466                if let Ok(exec_algorithm_id) = value.extract::<ExecAlgorithmId>() {
1467                    Ok(exec_algorithm_id)
1468                } else {
1469                    let value: String = value.extract()?;
1470                    ExecAlgorithmId::new_checked(&value).map_err(to_pyvalue_err)
1471                }
1472            })
1473            .transpose()?;
1474
1475        Ok(Self {
1476            exec_algorithm_id,
1477            log_events,
1478            log_commands,
1479        })
1480    }
1481
1482    #[getter]
1483    fn exec_algorithm_id(&self) -> Option<ExecAlgorithmId> {
1484        self.exec_algorithm_id
1485    }
1486
1487    #[getter]
1488    fn log_events(&self) -> bool {
1489        self.log_events
1490    }
1491
1492    #[getter]
1493    fn log_commands(&self) -> bool {
1494        self.log_commands
1495    }
1496}
1497
1498#[pyo3::pymethods]
1499#[pyo3_stub_gen::derive::gen_stub_pymethods]
1500impl ImportableExecutionAlgorithmConfig {
1501    /// Configuration for creating execution algorithms from importable paths.
1502    #[new]
1503    #[expect(clippy::needless_pass_by_value)]
1504    fn py_new(
1505        exec_algorithm_path: String,
1506        config_path: String,
1507        config: Py<PyDict>,
1508    ) -> PyResult<Self> {
1509        let json_config = Python::attach(|py| py_dict_to_json(config.bind(py)))?;
1510
1511        Ok(Self {
1512            exec_algorithm_path,
1513            config_path,
1514            config: json_config,
1515        })
1516    }
1517
1518    #[getter]
1519    fn exec_algorithm_path(&self) -> &String {
1520        &self.exec_algorithm_path
1521    }
1522
1523    #[getter]
1524    fn config_path(&self) -> &String {
1525        &self.config_path
1526    }
1527
1528    #[getter]
1529    fn config(&self, py: Python<'_>) -> PyResult<Py<PyDict>> {
1530        let py_dict = PyDict::new(py);
1531
1532        for (key, value) in &self.config {
1533            let json_str = serde_json::to_string(value).map_err(to_pyvalue_err)?;
1534            let py_value = PyModule::import(py, "json")?.call_method("loads", (json_str,), None)?;
1535            py_dict.set_item(key, py_value)?;
1536        }
1537        Ok(py_dict.unbind())
1538    }
1539}
1540
1541fn py_type_path(value: &Bound<'_, PyAny>) -> PyResult<String> {
1542    let value_type = value.get_type();
1543    let module: String = value_type.getattr("__module__")?.extract()?;
1544    let qualname: String = value_type.getattr("__qualname__")?.extract()?;
1545    Ok(format!("{module}:{qualname}"))
1546}
1547
1548fn py_config_to_json(config: &Bound<'_, PyAny>) -> PyResult<HashMap<String, serde_json::Value>> {
1549    let py = config.py();
1550    let config_dict = PyDict::new(py);
1551
1552    if let Ok(attributes) = config.getattr("__dict__")
1553        && let Ok(attributes) = attributes.cast::<PyDict>()
1554    {
1555        for (key, value) in attributes.iter() {
1556            config_dict.set_item(key, value)?;
1557        }
1558    }
1559
1560    for field in [
1561        "exec_algorithm_id",
1562        "actor_id",
1563        "log_events",
1564        "log_commands",
1565    ] {
1566        if let Ok(value) = config.getattr(field) {
1567            config_dict.set_item(field, value)?;
1568        }
1569    }
1570
1571    py_dict_to_json(&config_dict)
1572}
1573
1574fn py_dict_to_json(config: &Bound<'_, PyDict>) -> PyResult<HashMap<String, serde_json::Value>> {
1575    let py = config.py();
1576    let kwargs = PyDict::new(py);
1577    kwargs.set_item("default", py.eval(pyo3::ffi::c_str!("str"), None, None)?)?;
1578    let json_str: String = PyModule::import(py, "json")?
1579        .call_method("dumps", (config,), Some(&kwargs))?
1580        .extract()?;
1581
1582    let json_value: serde_json::Value = serde_json::from_str(&json_str).map_err(to_pyvalue_err)?;
1583
1584    if let serde_json::Value::Object(map) = json_value {
1585        Ok(map.into_iter().collect())
1586    } else {
1587        Err(to_pyvalue_err("Config must be a dictionary"))
1588    }
1589}
1590
1591#[cfg(test)]
1592mod tests {
1593    use std::{cell::RefCell, rc::Rc};
1594
1595    use nautilus_common::{
1596        cache::Cache,
1597        clock::{Clock, VirtualClock},
1598        logging::{
1599            arm_shutdown_on_error, disarm_shutdown_on_error, init_logging,
1600            take_shutdown_on_error_trigger,
1601        },
1602        messages::system::{
1603            QueueCondition, QueueState, QueueStateChanged, SocketState, SocketStateChanged,
1604        },
1605        msgbus::{
1606            MessageBus, MessagingSwitchboard, get_message_bus, switchboard::get_signal_topic,
1607        },
1608        runner::SystemChannel,
1609    };
1610    use nautilus_core::{UUID4, UnixNanos};
1611    use nautilus_model::{
1612        enums::{OrderSide, OrderType, TriggerType},
1613        events::order::spec::{OrderFillVoidedSpec, OrderFilledSpec},
1614        identifiers::{
1615            ClientId, ClientOrderId, InstrumentId, OrderListId, PositionId, StrategyId, TraderId,
1616            Venue,
1617        },
1618        instruments::{InstrumentAny, stubs::audusd_sim},
1619        orders::OrderTestBuilder,
1620        position::Position,
1621        types::{Price, Quantity},
1622    };
1623    use pyo3::{
1624        ffi::c_str,
1625        types::{PyWeakrefMethods, PyWeakrefReference},
1626    };
1627    use rstest::rstest;
1628    use ustr::Ustr;
1629
1630    use super::*;
1631
1632    #[rstest]
1633    fn test_python_submit_order_maps_emulation_refusal_to_value_error() {
1634        Python::initialize();
1635
1636        Python::attach(|py| {
1637            *get_message_bus().borrow_mut() = MessageBus::default();
1638
1639            let mut algorithm = PyExecutionAlgorithm::new(None);
1640            let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
1641            let cache = Rc::new(RefCell::new(Cache::default()));
1642            Component::register(&mut algorithm, TraderId::from("TRADER-001"), clock, cache)
1643                .unwrap();
1644
1645            let client_order_id = ClientOrderId::from("O-EMULATED-SUBMISSION");
1646            let emulated_order = OrderTestBuilder::new(OrderType::Limit)
1647                .instrument_id(InstrumentId::from("AUDUSD.SIM"))
1648                .client_order_id(client_order_id)
1649                .side(OrderSide::Buy)
1650                .price(Price::from("1.00000"))
1651                .quantity(Quantity::from("100"))
1652                .emulation_trigger(TriggerType::BidAsk)
1653                .build();
1654            let emulated_order =
1655                nautilus_model::python::orders::order_any_to_pyobject(py, emulated_order).unwrap();
1656
1657            let e = algorithm
1658                .py_submit_order(py, emulated_order, None, None)
1659                .unwrap_err();
1660            assert!(e.is_instance_of::<pyo3::exceptions::PyValueError>(py));
1661            let message = e.to_string();
1662            assert!(message.contains(client_order_id.as_str()), "{message}");
1663            assert!(message.contains("live emulation trigger"), "{message}");
1664
1665            let mut unregistered_algorithm = PyExecutionAlgorithm::new(None);
1666            let order = OrderTestBuilder::new(OrderType::Market)
1667                .instrument_id(InstrumentId::from("AUDUSD.SIM"))
1668                .side(OrderSide::Buy)
1669                .quantity(Quantity::from("100"))
1670                .build();
1671            let order = nautilus_model::python::orders::order_any_to_pyobject(py, order).unwrap();
1672            let e = unregistered_algorithm
1673                .py_submit_order(py, order, None, None)
1674                .unwrap_err();
1675            assert!(e.is_instance_of::<pyo3::exceptions::PyRuntimeError>(py));
1676        });
1677    }
1678
1679    fn sample_queue_state_changed() -> QueueStateChanged {
1680        QueueStateChanged::new(
1681            TraderId::from("TRADER-001"),
1682            SystemChannel::ExecCommands,
1683            QueueCondition::Backlogged,
1684            QueueState::Triggered,
1685            17,
1686            23,
1687            UUID4::from("00000000-0000-4000-8000-000000000001"),
1688            UnixNanos::from(1_700_000_000_000_000_001),
1689            UnixNanos::from(1_700_000_000_000_000_002),
1690        )
1691    }
1692
1693    fn sample_socket_state_changed() -> SocketStateChanged {
1694        SocketStateChanged::new(
1695            TraderId::from("TRADER-001"),
1696            ClientId::from("BINANCE"),
1697            Some(Venue::from("BINANCE")),
1698            Ustr::from("binance-futures-market-streams"),
1699            SocketState::Connected,
1700            UUID4::from("00000000-0000-4000-8000-000000000001"),
1701            UnixNanos::from(1_700_000_000_000_000_001),
1702            UnixNanos::from(1_700_000_000_000_000_002),
1703        )
1704    }
1705
1706    #[rstest]
1707    fn test_python_queue_state_dispatches_exact_event() {
1708        Python::initialize();
1709
1710        let tracker = Python::attach(|py| {
1711            py.run(
1712                c_str!(
1713                    r#"
1714class QueueStateTracker:
1715    def __init__(self):
1716        self.event = None
1717
1718    def on_queue_state(self, event):
1719        self.event = event
1720"#
1721                ),
1722                None,
1723                None,
1724            )
1725            .unwrap();
1726            py.eval(c_str!("QueueStateTracker()"), None, None)
1727                .unwrap()
1728                .unbind()
1729        });
1730        let mut algorithm = PyExecutionAlgorithm::new(None);
1731        Python::attach(|py| algorithm.set_python_instance(tracker.bind(py))).unwrap();
1732        let event = sample_queue_state_changed();
1733
1734        DataActor::on_queue_state(&mut algorithm, &event).unwrap();
1735
1736        let received = Python::attach(|py| {
1737            tracker
1738                .getattr(py, "event")
1739                .unwrap()
1740                .extract::<QueueStateChanged>(py)
1741                .unwrap()
1742        });
1743        assert_eq!(received, event);
1744    }
1745
1746    #[rstest]
1747    fn test_python_socket_state_dispatches_exact_event() {
1748        Python::initialize();
1749
1750        let tracker = Python::attach(|py| {
1751            py.run(
1752                c_str!(
1753                    r#"
1754class SocketStateTracker:
1755    def __init__(self):
1756        self.event = None
1757
1758    def on_socket_state(self, event):
1759        self.event = event
1760"#
1761                ),
1762                None,
1763                None,
1764            )
1765            .unwrap();
1766            py.eval(c_str!("SocketStateTracker()"), None, None)
1767                .unwrap()
1768                .unbind()
1769        });
1770        let mut algorithm = PyExecutionAlgorithm::new(None);
1771        Python::attach(|py| algorithm.set_python_instance(tracker.bind(py))).unwrap();
1772        let event = sample_socket_state_changed();
1773
1774        DataActor::on_socket_state(&mut algorithm, &event).unwrap();
1775
1776        let received = Python::attach(|py| {
1777            tracker
1778                .getattr(py, "event")
1779                .unwrap()
1780                .extract::<SocketStateChanged>(py)
1781                .unwrap()
1782        });
1783        assert_eq!(received, event);
1784    }
1785
1786    #[rstest]
1787    fn test_python_subscribe_and_unsubscribe_signal_update_msgbus() {
1788        *get_message_bus().borrow_mut() = MessageBus::default();
1789
1790        let mut algorithm = PyExecutionAlgorithm::new(None);
1791        let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
1792        let cache = Rc::new(RefCell::new(Cache::default()));
1793        Component::register(&mut algorithm, TraderId::from("TRADER-001"), clock, cache).unwrap();
1794
1795        algorithm.py_subscribe_signal("risk", Some(50)).unwrap();
1796
1797        let topic = get_signal_topic("risk");
1798        let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
1799        assert_eq!(subscriptions.len(), 1);
1800        assert_eq!(subscriptions[0].priority, 50);
1801
1802        algorithm.py_unsubscribe_signal("risk").unwrap();
1803
1804        let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
1805        assert!(subscriptions.is_empty());
1806    }
1807
1808    #[rstest]
1809    #[case(None)]
1810    #[case(Some(SystemChannel::ExecCommands))]
1811    fn test_python_subscribe_and_unsubscribe_queue_state_update_msgbus(
1812        #[case] channel: Option<SystemChannel>,
1813    ) {
1814        *get_message_bus().borrow_mut() = MessageBus::default();
1815
1816        let mut algorithm = PyExecutionAlgorithm::new(None);
1817        let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
1818        let cache = Rc::new(RefCell::new(Cache::default()));
1819        Component::register(&mut algorithm, TraderId::from("TRADER-001"), clock, cache).unwrap();
1820
1821        algorithm
1822            .py_subscribe_queue_state(channel, Some(50))
1823            .unwrap();
1824
1825        let topic = MessagingSwitchboard::queue_state_changed_topic(SystemChannel::ExecCommands);
1826        let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
1827        assert_eq!(subscriptions.len(), 1);
1828        assert_eq!(subscriptions[0].priority, 50);
1829        let unrelated = MessagingSwitchboard::queue_state_changed_topic(SystemChannel::DataEvents);
1830        assert_eq!(
1831            get_message_bus()
1832                .borrow_mut()
1833                .matching_subscriptions(unrelated)
1834                .len(),
1835            usize::from(channel.is_none())
1836        );
1837
1838        algorithm.py_unsubscribe_queue_state(channel).unwrap();
1839
1840        let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
1841        assert!(subscriptions.is_empty());
1842    }
1843
1844    #[rstest]
1845    #[case(None, None)]
1846    #[case(Some(ClientId::from("BINANCE")), None)]
1847    #[case(None, Some("binance-futures-market-streams"))]
1848    #[case(
1849        Some(ClientId::from("BINANCE")),
1850        Some("binance-futures-market-streams")
1851    )]
1852    fn test_python_subscribe_and_unsubscribe_socket_state_update_msgbus(
1853        #[case] client_id: Option<ClientId>,
1854        #[case] endpoint: Option<&str>,
1855    ) {
1856        *get_message_bus().borrow_mut() = MessageBus::default();
1857
1858        let mut algorithm = PyExecutionAlgorithm::new(None);
1859        let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
1860        let cache = Rc::new(RefCell::new(Cache::default()));
1861        Component::register(&mut algorithm, TraderId::from("TRADER-001"), clock, cache).unwrap();
1862
1863        algorithm
1864            .py_subscribe_socket_state(client_id, endpoint, Some(50))
1865            .unwrap();
1866
1867        let topic = MessagingSwitchboard::socket_state_changed_topic(
1868            ClientId::from("BINANCE"),
1869            "binance-futures-market-streams",
1870        );
1871        let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
1872        assert_eq!(subscriptions.len(), 1);
1873        assert_eq!(subscriptions[0].priority, 50);
1874        let unrelated =
1875            MessagingSwitchboard::socket_state_changed_topic(ClientId::from("BYBIT"), "orders");
1876        assert_eq!(
1877            get_message_bus()
1878                .borrow_mut()
1879                .matching_subscriptions(unrelated)
1880                .len(),
1881            usize::from(client_id.is_none() && endpoint.is_none())
1882        );
1883
1884        algorithm
1885            .py_unsubscribe_socket_state(client_id, endpoint)
1886            .unwrap();
1887
1888        let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
1889        assert!(subscriptions.is_empty());
1890    }
1891
1892    #[rstest]
1893    fn test_python_order_list_override_receives_resolved_orders_without_fanout() {
1894        Python::initialize();
1895
1896        let tracker = Python::attach(|py| {
1897            py.run(
1898                c_str!(
1899                    r#"
1900class OrderListTracker:
1901    def __init__(self):
1902        self.list_calls = 0
1903        self.list_ids = []
1904        self.resolved_ids = []
1905        self.order_ids = []
1906
1907    def on_order_list(self, order_list, orders):
1908        self.list_calls += 1
1909        self.list_ids = [str(value) for value in order_list.client_order_ids()]
1910        self.resolved_ids = [str(order.client_order_id) for order in orders]
1911
1912    def on_order(self, order):
1913        self.order_ids.append(str(order.client_order_id))
1914
1915    def observations(self):
1916        return self.list_calls, self.list_ids, self.resolved_ids, self.order_ids
1917"#
1918                ),
1919                None,
1920                None,
1921            )
1922            .unwrap();
1923            py.eval(c_str!("OrderListTracker()"), None, None)
1924                .unwrap()
1925                .unbind()
1926        });
1927        let mut algorithm = PyExecutionAlgorithm::new(None);
1928        Python::attach(|py| algorithm.set_python_instance(tracker.bind(py))).unwrap();
1929
1930        let instrument_id = InstrumentId::from("BTC/USDT.BINANCE");
1931        let strategy_id = StrategyId::from("STRAT-LIST-OVERRIDE");
1932        let first = OrderTestBuilder::new(OrderType::Market)
1933            .strategy_id(strategy_id)
1934            .instrument_id(instrument_id)
1935            .client_order_id(ClientOrderId::from("O-LIST-OVERRIDE-001"))
1936            .quantity(Quantity::from("1.0"))
1937            .build();
1938        let second = OrderTestBuilder::new(OrderType::Market)
1939            .strategy_id(strategy_id)
1940            .instrument_id(instrument_id)
1941            .client_order_id(ClientOrderId::from("O-LIST-OVERRIDE-002"))
1942            .quantity(Quantity::from("2.0"))
1943            .build();
1944        let order_list = OrderList::new(
1945            OrderListId::from("OL-OVERRIDE-001"),
1946            instrument_id,
1947            strategy_id,
1948            vec![first.client_order_id(), second.client_order_id()],
1949            0.into(),
1950        );
1951
1952        ExecutionAlgorithm::on_order_list(&mut algorithm, order_list, vec![first, second]).unwrap();
1953
1954        let observations = Python::attach(|py| {
1955            tracker
1956                .call_method0(py, "observations")
1957                .unwrap()
1958                .extract::<(usize, Vec<String>, Vec<String>, Vec<String>)>(py)
1959                .unwrap()
1960        });
1961        assert_eq!(
1962            observations,
1963            (
1964                1,
1965                vec![
1966                    "O-LIST-OVERRIDE-001".to_string(),
1967                    "O-LIST-OVERRIDE-002".to_string(),
1968                ],
1969                vec![
1970                    "O-LIST-OVERRIDE-001".to_string(),
1971                    "O-LIST-OVERRIDE-002".to_string(),
1972                ],
1973                Vec::new(),
1974            ),
1975        );
1976    }
1977
1978    #[rstest]
1979    fn test_python_self_is_weak() {
1980        Python::initialize();
1981
1982        Python::attach(|py| {
1983            let instance = py
1984                .get_type::<PyExecutionAlgorithm>()
1985                .call0()
1986                .expect("ExecutionAlgorithm should construct");
1987            let weakref = PyWeakrefReference::new(&instance)
1988                .expect("ExecutionAlgorithm should be weak-referenceable");
1989            assert!(weakref.upgrade().is_some());
1990
1991            drop(instance);
1992
1993            // A strong `py_self` would form an untraceable Rust-Python cycle and keep this alive
1994            assert!(
1995                weakref.upgrade().is_none(),
1996                "an unregistered ExecutionAlgorithm must be collected once its last Python owner is dropped",
1997            );
1998        });
1999    }
2000
2001    #[rstest]
2002    #[case("on_order_initialized")]
2003    #[case("on_order_denied")]
2004    #[case("on_order_emulated")]
2005    #[case("on_order_released")]
2006    #[case("on_order_submitted")]
2007    #[case("on_order_rejected")]
2008    #[case("on_order_accepted")]
2009    #[case("on_order_canceled")]
2010    #[case("on_order_expired")]
2011    #[case("on_order_triggered")]
2012    #[case("on_order_pending_update")]
2013    #[case("on_order_pending_cancel")]
2014    #[case("on_order_modify_rejected")]
2015    #[case("on_order_cancel_rejected")]
2016    #[case("on_order_updated")]
2017    #[case("on_order_filled")]
2018    #[case("on_order_fill_voided")]
2019    #[case("on_order_event")]
2020    #[case("on_position_opened")]
2021    #[case("on_position_changed")]
2022    #[case("on_position_closed")]
2023    #[case("on_position_event")]
2024    fn test_python_callback_errors(
2025        #[case] method: &str,
2026        #[values(false, true)] raises: bool,
2027        #[values(false, true)] shutdown: bool,
2028    ) {
2029        Python::initialize();
2030        let _guard = init_logging(
2031            TraderId::from("TRADER-001"),
2032            UUID4::new(),
2033            Default::default(),
2034            Default::default(),
2035        )
2036        .unwrap();
2037        Python::attach(|py| {
2038            let fill = OrderFilledSpec::builder()
2039                .position_id(PositionId::from("P-001"))
2040                .build();
2041            let position = Position::new(&InstrumentAny::CurrencyPair(audusd_sim()), fill.clone());
2042            let tracker = create_callback_algorithm(py, raises);
2043            let mut algorithm = PyExecutionAlgorithm::new(None);
2044            algorithm.set_python_instance(tracker.bind(py)).unwrap();
2045            arm_shutdown_on_error(shutdown);
2046
2047            match method {
2048                "on_order_initialized" => ExecutionAlgorithm::on_order_initialized(
2049                    &mut algorithm,
2050                    OrderInitialized::default(),
2051                ),
2052                "on_order_denied" => {
2053                    ExecutionAlgorithm::on_order_denied(&mut algorithm, OrderDenied::default());
2054                }
2055                "on_order_emulated" => {
2056                    ExecutionAlgorithm::on_order_emulated(&mut algorithm, OrderEmulated::default());
2057                }
2058                "on_order_released" => {
2059                    ExecutionAlgorithm::on_order_released(&mut algorithm, OrderReleased::default());
2060                }
2061                "on_order_submitted" => ExecutionAlgorithm::on_order_submitted(
2062                    &mut algorithm,
2063                    OrderSubmitted::default(),
2064                ),
2065                "on_order_rejected" => {
2066                    ExecutionAlgorithm::on_order_rejected(&mut algorithm, OrderRejected::default());
2067                }
2068                "on_order_accepted" => {
2069                    ExecutionAlgorithm::on_order_accepted(&mut algorithm, OrderAccepted::default());
2070                }
2071                "on_order_canceled" => ExecutionAlgorithm::on_algo_order_canceled(
2072                    &mut algorithm,
2073                    OrderCanceled::default(),
2074                ),
2075                "on_order_expired" => {
2076                    ExecutionAlgorithm::on_order_expired(&mut algorithm, OrderExpired::default());
2077                }
2078                "on_order_triggered" => ExecutionAlgorithm::on_order_triggered(
2079                    &mut algorithm,
2080                    OrderTriggered::default(),
2081                ),
2082                "on_order_pending_update" => ExecutionAlgorithm::on_order_pending_update(
2083                    &mut algorithm,
2084                    OrderPendingUpdate::default(),
2085                ),
2086                "on_order_pending_cancel" => ExecutionAlgorithm::on_order_pending_cancel(
2087                    &mut algorithm,
2088                    OrderPendingCancel::default(),
2089                ),
2090                "on_order_modify_rejected" => ExecutionAlgorithm::on_order_modify_rejected(
2091                    &mut algorithm,
2092                    OrderModifyRejected::default(),
2093                ),
2094                "on_order_cancel_rejected" => ExecutionAlgorithm::on_order_cancel_rejected(
2095                    &mut algorithm,
2096                    OrderCancelRejected::default(),
2097                ),
2098                "on_order_updated" => {
2099                    ExecutionAlgorithm::on_order_updated(&mut algorithm, OrderUpdated::default());
2100                }
2101                "on_order_filled" => ExecutionAlgorithm::on_algo_order_filled(
2102                    &mut algorithm,
2103                    OrderFilledSpec::builder().build(),
2104                ),
2105                "on_order_fill_voided" => ExecutionAlgorithm::on_order_fill_voided(
2106                    &mut algorithm,
2107                    &OrderFillVoidedSpec::builder().build(),
2108                ),
2109                "on_order_event" => ExecutionAlgorithm::on_order_event(
2110                    &mut algorithm,
2111                    OrderEventAny::Accepted(OrderAccepted::default()),
2112                ),
2113                "on_position_opened" => ExecutionAlgorithm::on_position_opened(
2114                    &mut algorithm,
2115                    PositionOpened::create(&position, &fill, UUID4::new(), UnixNanos::default()),
2116                ),
2117                "on_position_changed" => ExecutionAlgorithm::on_position_changed(
2118                    &mut algorithm,
2119                    PositionChanged::create(&position, &fill, UUID4::new(), UnixNanos::default()),
2120                ),
2121                "on_position_closed" => ExecutionAlgorithm::on_position_closed(
2122                    &mut algorithm,
2123                    PositionClosed::create(&position, &fill, UUID4::new(), UnixNanos::default()),
2124                ),
2125                "on_position_event" => ExecutionAlgorithm::on_position_event(
2126                    &mut algorithm,
2127                    PositionEvent::PositionOpened(PositionOpened::create(
2128                        &position,
2129                        &fill,
2130                        UUID4::new(),
2131                        UnixNanos::default(),
2132                    )),
2133                ),
2134                _ => unreachable!(),
2135            }
2136
2137            let trigger = take_shutdown_on_error_trigger();
2138            disarm_shutdown_on_error();
2139            let calls = tracker
2140                .getattr(py, "calls")
2141                .unwrap()
2142                .extract::<Vec<String>>(py)
2143                .unwrap();
2144
2145            assert_eq!(calls, [method]);
2146
2147            if raises && shutdown {
2148                let trigger = trigger.expect("Python callback failure must request shutdown");
2149                assert_eq!(
2150                    trigger.component.as_str(),
2151                    algorithm.inner().core.actor.actor_id.as_str()
2152                );
2153                assert!(
2154                    trigger
2155                        .message
2156                        .contains(&format!("Python {method} failed:"))
2157                );
2158                assert!(trigger.message.contains("in callback"));
2159                assert!(
2160                    trigger
2161                        .message
2162                        .contains("RuntimeError: algorithm callback failure")
2163                );
2164            } else {
2165                assert_eq!(trigger, None);
2166            }
2167        });
2168    }
2169
2170    #[rstest]
2171    fn test_python_algorithm_callback_exception_preserves_traceback() {
2172        Python::initialize();
2173        Python::attach(|py| {
2174            let tracker = create_callback_algorithm(py, true);
2175            let mut algorithm = PyExecutionAlgorithm::new(None);
2176            algorithm.set_python_instance(tracker.bind(py)).unwrap();
2177
2178            let event = TimeEvent::new(
2179                "ALERT".into(),
2180                UUID4::new(),
2181                UnixNanos::from(10),
2182                UnixNanos::from(11),
2183            );
2184            let error = DataActor::on_time_event(&mut algorithm, &event).unwrap_err();
2185            let message = error.to_string();
2186            let calls = tracker
2187                .getattr(py, "calls")
2188                .unwrap()
2189                .extract::<Vec<String>>(py)
2190                .unwrap();
2191
2192            assert_eq!(calls, ["on_time_event"]);
2193            assert!(message.contains("Python on_time_event failed:"));
2194            assert!(message.contains("in callback"));
2195            assert!(message.contains("RuntimeError: algorithm callback failure"));
2196        });
2197    }
2198
2199    fn create_callback_algorithm(py: Python<'_>, raises: bool) -> Py<PyAny> {
2200        let module = PyModule::from_code(
2201            py,
2202            c_str!(
2203                r#"
2204class CallbackAlgorithm:
2205    def __init__(self, raises):
2206        self.raises = raises
2207        self.calls = []
2208
2209    def __getattr__(self, name):
2210        def callback(event):
2211            self.calls.append(name)
2212            if self.raises:
2213                raise RuntimeError("algorithm callback failure")
2214        return callback
2215"#
2216            ),
2217            c_str!("algorithm_callback.py"),
2218            c_str!("algorithm_callback"),
2219        )
2220        .unwrap();
2221        module
2222            .getattr("CallbackAlgorithm")
2223            .unwrap()
2224            .call1((raises,))
2225            .unwrap()
2226            .unbind()
2227    }
2228}