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