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