1use 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
62pub struct PyExecutionAlgorithmInner {
64 core: ExecutionAlgorithmCore,
65 py_self: Option<Py<PyWeakrefReference>>,
66 config: Option<Py<PyAny>>,
67 logger: PyLogger,
68}
69
70impl PyExecutionAlgorithmInner {
71 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#[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 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 unsafe { &mut *self.inner.get() }
131 }
132}
133
134impl PyExecutionAlgorithm {
135 #[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 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 pub fn set_config(&mut self, config: Option<Py<PyAny>>) {
173 self.inner_mut().config = config;
174 }
175
176 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 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 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 #[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 #[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 #[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 #[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 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 #[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 #[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 assert!(
1818 weakref.upgrade().is_none(),
1819 "an unregistered ExecutionAlgorithm must be collected once its last Python owner is dropped",
1820 );
1821 });
1822 }
1823}