1use std::{
19 any::Any,
20 cell::{RefCell, UnsafeCell},
21 collections::HashMap,
22 fmt::Debug,
23 num::NonZeroUsize,
24 ops::{Deref, DerefMut},
25 rc::Rc,
26};
27
28use indexmap::IndexMap;
29use jiff::Timestamp;
30use nautilus_common::{
31 actor::{
32 Actor, DataActor, DataActorNative,
33 data_actor::DataActorCore,
34 registry::{try_get_actor_unchecked, with_actor_registry},
35 },
36 cache::Cache,
37 clock::Clock,
38 component::{Component, with_component_registry},
39 enums::ComponentState,
40 messages::system::{QueueStateChanged, SocketStateChanged},
41 python::{
42 cache::PyCache,
43 clock::PyClock,
44 config_error_to_pyvalue_err,
45 indicators::{registered_python_indicators, wrap_python_indicator},
46 logging::{PyLogger, format_exception},
47 order_factory::PyOrderFactory,
48 wrappers::{get_python_message_bus, retain_python_wrapper},
49 },
50 runner::SystemChannel,
51 signal::Signal,
52 timer::{TimeEvent, TimeEventCallback},
53};
54use nautilus_core::{
55 Params, UnixNanos,
56 correctness::{CorrectnessResult, CorrectnessResultExt, FAILED},
57 from_pydict,
58 python::{to_pyruntime_err, to_pyvalue_err, upgrade_py_weakref},
59};
60use nautilus_model::{
61 data::{
62 Bar, BarType, CustomData, DataType, FundingRateUpdate, IndexPriceUpdate, InstrumentStatus,
63 MarkPriceUpdate, OrderBookDelta, OrderBookDeltas, OrderBookDepth, QuoteTick, TradeTick,
64 close::InstrumentClose,
65 option_chain::{OptionChainSlice, OptionGreeks},
66 },
67 enums::{BookType, OmsType, OrderSide, PositionSide, TimeInForce},
68 events::{
69 OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
70 OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled, OrderInitialized,
71 OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
72 OrderSubmitted, OrderTriggered, OrderUpdated, PositionChanged, PositionClosed,
73 PositionEvent, PositionOpened,
74 },
75 identifiers::{
76 AccountId, ActorId, ClientId, ClientOrderId, InstrumentId, OptionSeriesId, PositionId,
77 StrategyId, TraderId, UNASSIGNED_ORDER_ID_TAG, Venue, normalize_order_id_tag,
78 },
79 instruments::{InstrumentAny, SyntheticInstrument},
80 orderbook::OrderBook,
81 orders::{Order, OrderAny},
82 position::Position,
83 python::{
84 data::option_chain::PyStrikeRange, events::order::order_event_to_pyobject,
85 instruments::instrument_any_to_pyobject, orders::pyobject_to_order_any,
86 },
87 types::{Price, Quantity},
88};
89use nautilus_portfolio::{portfolio::Portfolio, python::PyPortfolio};
90use pyo3::{
91 IntoPyObjectExt,
92 prelude::*,
93 types::{PyBytes, PyDict, PyList, PyWeakrefReference},
94};
95use ustr::Ustr;
96
97use crate::strategy::{
98 BatchModifyOrder, ImportableStrategyConfig, Strategy, StrategyConfig, StrategyCore,
99 StrategyNative, route_time_event,
100};
101
102#[pyo3::pymethods]
103#[pyo3_stub_gen::derive::gen_stub_pymethods]
104impl StrategyConfig {
105 #[new]
107 #[pyo3(signature = (
108 strategy_id=None,
109 order_id_tag=None,
110 oms_type=None,
111 external_order_instrument_ids=None,
112 manage_contingent_orders=false,
113 manage_gtd_expiry=false,
114 manage_stop=false,
115 market_exit_interval_ms=100,
116 market_exit_max_attempts=100,
117 market_exit_time_in_force=TimeInForce::Gtc,
118 market_exit_reduce_only=true,
119 use_uuid_client_order_ids=false,
120 use_hyphens_in_client_order_ids=true,
121 log_events=true,
122 log_commands=true,
123 log_rejected_due_post_only_as_warning=true,
124 **_kwargs
125 ))]
126 #[expect(
127 clippy::fn_params_excessive_bools,
128 clippy::too_many_arguments,
129 reason = "constructor mirrors the existing Python keyword API"
130 )]
131 fn py_new(
132 strategy_id: Option<StrategyId>,
133 order_id_tag: Option<String>,
134 oms_type: Option<OmsType>,
135 external_order_instrument_ids: Option<Vec<InstrumentId>>,
136 manage_contingent_orders: bool,
137 manage_gtd_expiry: bool,
138 manage_stop: bool,
139 market_exit_interval_ms: u64,
140 market_exit_max_attempts: u64,
141 market_exit_time_in_force: TimeInForce,
142 market_exit_reduce_only: bool,
143 use_uuid_client_order_ids: bool,
144 use_hyphens_in_client_order_ids: bool,
145 log_events: bool,
146 log_commands: bool,
147 log_rejected_due_post_only_as_warning: bool,
148 _kwargs: Option<&Bound<'_, PyDict>>,
149 ) -> PyResult<Self> {
150 let config = Self {
151 strategy_id,
152 order_id_tag,
153 use_uuid_client_order_ids,
154 use_hyphens_in_client_order_ids,
155 oms_type,
156 external_order_instrument_ids,
157 manage_contingent_orders,
158 manage_gtd_expiry,
159 manage_stop,
160 market_exit_interval_ms,
161 market_exit_max_attempts,
162 market_exit_time_in_force,
163 market_exit_reduce_only,
164 log_events,
165 log_commands,
166 log_rejected_due_post_only_as_warning,
167 };
168 config.validate().map_err(config_error_to_pyvalue_err)?;
169 Ok(config)
170 }
171
172 #[getter]
173 fn strategy_id(&self) -> Option<StrategyId> {
174 self.strategy_id
175 }
176
177 #[getter]
178 fn order_id_tag(&self) -> Option<&String> {
179 self.order_id_tag.as_ref()
180 }
181
182 #[getter]
183 fn oms_type(&self) -> Option<OmsType> {
184 self.oms_type
185 }
186
187 #[getter]
188 fn external_order_instrument_ids(&self) -> Option<Vec<InstrumentId>> {
189 self.external_order_instrument_ids.clone()
190 }
191
192 #[getter]
193 fn manage_contingent_orders(&self) -> bool {
194 self.manage_contingent_orders
195 }
196
197 #[getter]
198 fn manage_gtd_expiry(&self) -> bool {
199 self.manage_gtd_expiry
200 }
201
202 #[getter]
203 fn manage_stop(&self) -> bool {
204 self.manage_stop
205 }
206
207 #[getter]
208 fn market_exit_interval_ms(&self) -> u64 {
209 self.market_exit_interval_ms
210 }
211
212 #[getter]
213 fn market_exit_max_attempts(&self) -> u64 {
214 self.market_exit_max_attempts
215 }
216
217 #[getter]
218 fn market_exit_time_in_force(&self) -> TimeInForce {
219 self.market_exit_time_in_force
220 }
221
222 #[getter]
223 fn market_exit_reduce_only(&self) -> bool {
224 self.market_exit_reduce_only
225 }
226
227 #[getter]
228 fn use_uuid_client_order_ids(&self) -> bool {
229 self.use_uuid_client_order_ids
230 }
231
232 #[getter]
233 fn use_hyphens_in_client_order_ids(&self) -> bool {
234 self.use_hyphens_in_client_order_ids
235 }
236
237 #[getter]
238 fn log_events(&self) -> bool {
239 self.log_events
240 }
241
242 #[getter]
243 fn log_commands(&self) -> bool {
244 self.log_commands
245 }
246
247 #[getter]
248 fn log_rejected_due_post_only_as_warning(&self) -> bool {
249 self.log_rejected_due_post_only_as_warning
250 }
251}
252
253#[pyo3::pymethods]
254#[pyo3_stub_gen::derive::gen_stub_pymethods]
255impl ImportableStrategyConfig {
256 #[new]
258 #[expect(clippy::needless_pass_by_value)]
259 fn py_new(strategy_path: String, config_path: String, config: Py<PyDict>) -> PyResult<Self> {
260 let json_config = Python::attach(|py| -> PyResult<HashMap<String, serde_json::Value>> {
261 let kwargs = PyDict::new(py);
262 kwargs.set_item("default", py.eval(pyo3::ffi::c_str!("str"), None, None)?)?;
263 let json_str: String = PyModule::import(py, "json")?
264 .call_method("dumps", (config.bind(py),), Some(&kwargs))?
265 .extract()?;
266
267 let json_value: serde_json::Value =
268 serde_json::from_str(&json_str).map_err(to_pyvalue_err)?;
269
270 if let serde_json::Value::Object(map) = json_value {
271 Ok(map.into_iter().collect())
272 } else {
273 Err(to_pyvalue_err("Config must be a dictionary"))
274 }
275 })?;
276
277 Ok(Self {
278 strategy_path,
279 config_path,
280 config: json_config,
281 })
282 }
283
284 #[getter]
285 fn strategy_path(&self) -> &String {
286 &self.strategy_path
287 }
288
289 #[getter]
290 fn config_path(&self) -> &String {
291 &self.config_path
292 }
293
294 #[getter]
295 fn config(&self, py: Python<'_>) -> PyResult<Py<PyDict>> {
296 let py_dict = PyDict::new(py);
297
298 for (key, value) in &self.config {
299 let json_str = serde_json::to_string(value).map_err(to_pyvalue_err)?;
300 let py_value = PyModule::import(py, "json")?.call_method("loads", (json_str,), None)?;
301 py_dict.set_item(key, py_value)?;
302 }
303 Ok(py_dict.unbind())
304 }
305}
306
307pub struct PyStrategyInner {
309 core: StrategyCore,
310 py_self: Option<Py<PyWeakrefReference>>,
311 config: Option<Py<PyAny>>,
312 clock: PyClock,
313 logger: PyLogger,
314}
315
316impl Debug for PyStrategyInner {
317 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
318 f.debug_struct(stringify!(PyStrategyInner))
319 .field("core", &self.core)
320 .field(
321 "py_self",
322 &self.py_self.as_ref().map(|_| "<Py<PyWeakrefReference>>"),
323 )
324 .field("config", &self.config.as_ref().map(|_| "<Py<PyAny>>"))
325 .field("clock", &self.clock)
326 .field("logger", &self.logger)
327 .finish()
328 }
329}
330
331#[expect(
332 clippy::needless_pass_by_ref_mut,
333 reason = "dispatch methods share receiver shape with mutable DataActor hooks"
334)]
335impl PyStrategyInner {
336 fn dispatch_on_start(&self) -> PyResult<()> {
337 if let Some(py_self) = self.python_instance()? {
338 Python::attach(|py| py_self.call_method0(py, "on_start"))?;
339 }
340 Ok(())
341 }
342
343 fn dispatch_on_stop(&self) -> PyResult<()> {
344 if let Some(py_self) = self.python_instance()? {
345 Python::attach(|py| py_self.call_method0(py, "on_stop"))?;
346 }
347 Ok(())
348 }
349
350 fn dispatch_on_resume(&self) -> PyResult<()> {
351 if let Some(py_self) = self.python_instance()? {
352 Python::attach(|py| py_self.call_method0(py, "on_resume"))?;
353 }
354 Ok(())
355 }
356
357 fn dispatch_on_reset(&self) -> PyResult<()> {
358 if let Some(py_self) = self.python_instance()? {
359 Python::attach(|py| py_self.call_method0(py, "on_reset"))?;
360 }
361 Ok(())
362 }
363
364 fn dispatch_on_dispose(&self) -> PyResult<()> {
365 if let Some(py_self) = self.python_instance()? {
366 Python::attach(|py| py_self.call_method0(py, "on_dispose"))?;
367 }
368 Ok(())
369 }
370
371 fn dispatch_on_degrade(&self) -> PyResult<()> {
372 if let Some(py_self) = self.python_instance()? {
373 Python::attach(|py| py_self.call_method0(py, "on_degrade"))?;
374 }
375 Ok(())
376 }
377
378 fn dispatch_on_fault(&self) -> PyResult<()> {
379 if let Some(py_self) = self.python_instance()? {
380 Python::attach(|py| py_self.call_method0(py, "on_fault"))?;
381 }
382 Ok(())
383 }
384
385 fn dispatch_on_save(&self) -> PyResult<IndexMap<String, Vec<u8>>> {
386 if let Some(py_self) = self.python_instance()? {
387 Python::attach(|py| {
388 let py_state = py_self.call_method0(py, "on_save")?;
389 let py_state: &Bound<'_, PyDict> = py_state.cast_bound::<PyDict>(py)?;
390 pydict_to_state(py_state)
391 })
392 } else {
393 Ok(IndexMap::new())
394 }
395 }
396
397 fn dispatch_on_load(&self, state: &IndexMap<String, Vec<u8>>) -> PyResult<()> {
398 if let Some(py_self) = self.python_instance()? {
399 Python::attach(|py| -> PyResult<()> {
400 let py_state = state_to_pydict(py, state)?;
401 py_self.call_method1(py, "on_load", (py_state,))?;
402 Ok(())
403 })?;
404 }
405 Ok(())
406 }
407
408 fn dispatch_on_market_exit(&self) -> PyResult<()> {
409 if let Some(py_self) = self.python_instance()? {
410 Python::attach(|py| py_self.call_method0(py, "on_market_exit"))?;
411 }
412 Ok(())
413 }
414
415 fn dispatch_post_market_exit(&self) -> PyResult<()> {
416 if let Some(py_self) = self.python_instance()? {
417 Python::attach(|py| py_self.call_method0(py, "post_market_exit"))?;
418 }
419 Ok(())
420 }
421
422 fn dispatch_on_time_event(&self, event: &TimeEvent) -> PyResult<()> {
423 if let Some(py_self) = self.python_instance()? {
424 Python::attach(|py| {
425 py_self.call_method1(py, "on_time_event", (event.clone().into_py_any(py)?,))
426 })?;
427 }
428 Ok(())
429 }
430
431 fn dispatch_on_order_initialized(&self, event: OrderInitialized) -> PyResult<()> {
432 if let Some(py_self) = self.python_instance()? {
433 Python::attach(|py| {
434 py_self.call_method1(py, "on_order_initialized", (event.into_py_any(py)?,))
435 })?;
436 }
437 Ok(())
438 }
439
440 fn dispatch_on_order_event(&self, event: OrderEventAny) -> PyResult<()> {
441 if let Some(py_self) = self.python_instance()? {
442 Python::attach(|py| {
443 let py_event = order_event_to_pyobject(py, event)?;
444 py_self.call_method1(py, "on_order_event", (py_event,))
445 })?;
446 }
447 Ok(())
448 }
449
450 fn dispatch_on_order_denied(&self, event: OrderDenied) -> PyResult<()> {
451 if let Some(py_self) = self.python_instance()? {
452 Python::attach(|py| {
453 py_self.call_method1(py, "on_order_denied", (event.into_py_any(py)?,))
454 })?;
455 }
456 Ok(())
457 }
458
459 fn dispatch_on_order_emulated(&self, event: OrderEmulated) -> PyResult<()> {
460 if let Some(py_self) = self.python_instance()? {
461 Python::attach(|py| {
462 py_self.call_method1(py, "on_order_emulated", (event.into_py_any(py)?,))
463 })?;
464 }
465 Ok(())
466 }
467
468 fn dispatch_on_order_released(&self, event: OrderReleased) -> PyResult<()> {
469 if let Some(py_self) = self.python_instance()? {
470 Python::attach(|py| {
471 py_self.call_method1(py, "on_order_released", (event.into_py_any(py)?,))
472 })?;
473 }
474 Ok(())
475 }
476
477 fn dispatch_on_order_submitted(&self, event: OrderSubmitted) -> PyResult<()> {
478 if let Some(py_self) = self.python_instance()? {
479 Python::attach(|py| {
480 py_self.call_method1(py, "on_order_submitted", (event.into_py_any(py)?,))
481 })?;
482 }
483 Ok(())
484 }
485
486 fn dispatch_on_order_rejected(&self, event: OrderRejected) -> PyResult<()> {
487 if let Some(py_self) = self.python_instance()? {
488 Python::attach(|py| {
489 py_self.call_method1(py, "on_order_rejected", (event.into_py_any(py)?,))
490 })?;
491 }
492 Ok(())
493 }
494
495 fn dispatch_on_order_accepted(&self, event: OrderAccepted) -> PyResult<()> {
496 if let Some(py_self) = self.python_instance()? {
497 Python::attach(|py| {
498 py_self.call_method1(py, "on_order_accepted", (event.into_py_any(py)?,))
499 })?;
500 }
501 Ok(())
502 }
503
504 fn dispatch_on_order_expired(&self, event: OrderExpired) -> PyResult<()> {
505 if let Some(py_self) = self.python_instance()? {
506 Python::attach(|py| {
507 py_self.call_method1(py, "on_order_expired", (event.into_py_any(py)?,))
508 })?;
509 }
510 Ok(())
511 }
512
513 fn dispatch_on_order_triggered(&self, event: OrderTriggered) -> PyResult<()> {
514 if let Some(py_self) = self.python_instance()? {
515 Python::attach(|py| {
516 py_self.call_method1(py, "on_order_triggered", (event.into_py_any(py)?,))
517 })?;
518 }
519 Ok(())
520 }
521
522 fn dispatch_on_order_pending_update(&self, event: OrderPendingUpdate) -> PyResult<()> {
523 if let Some(py_self) = self.python_instance()? {
524 Python::attach(|py| {
525 py_self.call_method1(py, "on_order_pending_update", (event.into_py_any(py)?,))
526 })?;
527 }
528 Ok(())
529 }
530
531 fn dispatch_on_order_pending_cancel(&self, event: OrderPendingCancel) -> PyResult<()> {
532 if let Some(py_self) = self.python_instance()? {
533 Python::attach(|py| {
534 py_self.call_method1(py, "on_order_pending_cancel", (event.into_py_any(py)?,))
535 })?;
536 }
537 Ok(())
538 }
539
540 fn dispatch_on_order_modify_rejected(&self, event: OrderModifyRejected) -> PyResult<()> {
541 if let Some(py_self) = self.python_instance()? {
542 Python::attach(|py| {
543 py_self.call_method1(py, "on_order_modify_rejected", (event.into_py_any(py)?,))
544 })?;
545 }
546 Ok(())
547 }
548
549 fn dispatch_on_order_cancel_rejected(&self, event: OrderCancelRejected) -> PyResult<()> {
550 if let Some(py_self) = self.python_instance()? {
551 Python::attach(|py| {
552 py_self.call_method1(py, "on_order_cancel_rejected", (event.into_py_any(py)?,))
553 })?;
554 }
555 Ok(())
556 }
557
558 fn dispatch_on_order_updated(&self, event: &OrderUpdated) -> PyResult<()> {
559 if let Some(py_self) = self.python_instance()? {
560 Python::attach(|py| {
561 py_self.call_method1(py, "on_order_updated", ((*event).into_py_any(py)?,))
562 })?;
563 }
564 Ok(())
565 }
566
567 fn dispatch_on_order_canceled(&self, event: OrderCanceled) -> PyResult<()> {
568 if let Some(py_self) = self.python_instance()? {
569 Python::attach(|py| {
570 py_self.call_method1(py, "on_order_canceled", (event.into_py_any(py)?,))
571 })?;
572 }
573 Ok(())
574 }
575
576 fn dispatch_on_order_filled(&self, event: &OrderFilled) -> PyResult<()> {
577 if let Some(py_self) = self.python_instance()? {
578 Python::attach(|py| {
579 py_self.call_method1(py, "on_order_filled", (event.clone().into_py_any(py)?,))
580 })?;
581 }
582 Ok(())
583 }
584
585 fn dispatch_on_order_fill_voided(&self, event: &OrderFillVoided) -> PyResult<()> {
586 if let Some(py_self) = self.python_instance()? {
587 Python::attach(|py| {
588 py_self.call_method1(
589 py,
590 "on_order_fill_voided",
591 (event.clone().into_py_any(py)?,),
592 )
593 })?;
594 }
595 Ok(())
596 }
597
598 fn dispatch_on_position_opened(&self, event: PositionOpened) -> PyResult<()> {
599 if let Some(py_self) = self.python_instance()? {
600 Python::attach(|py| {
601 py_self.call_method1(py, "on_position_opened", (event.into_py_any(py)?,))
602 })?;
603 }
604 Ok(())
605 }
606
607 fn dispatch_on_position_event(&self, event: PositionEvent) -> PyResult<()> {
608 if let Some(py_self) = self.python_instance()? {
609 Python::attach(|py| {
610 let py_event = match event {
611 PositionEvent::PositionOpened(event) => event.into_py_any(py)?,
612 PositionEvent::PositionChanged(event) => event.into_py_any(py)?,
613 PositionEvent::PositionClosed(event) => event.into_py_any(py)?,
614 PositionEvent::PositionAdjusted(event) => event.into_py_any(py)?,
615 };
616 py_self.call_method1(py, "on_position_event", (py_event,))
617 })?;
618 }
619 Ok(())
620 }
621
622 fn dispatch_on_position_changed(&self, event: PositionChanged) -> PyResult<()> {
623 if let Some(py_self) = self.python_instance()? {
624 Python::attach(|py| {
625 py_self.call_method1(py, "on_position_changed", (event.into_py_any(py)?,))
626 })?;
627 }
628 Ok(())
629 }
630
631 fn dispatch_on_position_closed(&self, event: PositionClosed) -> PyResult<()> {
632 if let Some(py_self) = self.python_instance()? {
633 Python::attach(|py| {
634 py_self.call_method1(py, "on_position_closed", (event.into_py_any(py)?,))
635 })?;
636 }
637 Ok(())
638 }
639
640 fn dispatch_on_data(&mut self, data: Py<PyAny>) -> PyResult<()> {
641 if let Some(py_self) = self.python_instance()? {
642 Python::attach(|py| py_self.call_method1(py, "on_data", (data,)))?;
643 }
644 Ok(())
645 }
646
647 fn dispatch_on_signal(&mut self, signal: &Signal) -> PyResult<()> {
648 if let Some(py_self) = self.python_instance()? {
649 Python::attach(|py| {
650 py_self.call_method1(py, "on_signal", (signal.clone().into_py_any(py)?,))
651 })?;
652 }
653 Ok(())
654 }
655
656 fn dispatch_on_queue_state(&mut self, event: &QueueStateChanged) -> PyResult<()> {
657 if let Some(py_self) = self.python_instance()? {
658 Python::attach(|py| {
659 py_self.call_method1(py, "on_queue_state", (event.clone().into_py_any(py)?,))
660 })?;
661 }
662 Ok(())
663 }
664
665 fn dispatch_on_socket_state(&mut self, event: &SocketStateChanged) -> PyResult<()> {
666 if let Some(py_self) = self.python_instance()? {
667 Python::attach(|py| {
668 py_self.call_method1(py, "on_socket_state", (event.clone().into_py_any(py)?,))
669 })?;
670 }
671 Ok(())
672 }
673
674 fn dispatch_on_instrument(&mut self, instrument: Py<PyAny>) -> PyResult<()> {
675 if let Some(py_self) = self.python_instance()? {
676 Python::attach(|py| py_self.call_method1(py, "on_instrument", (instrument,)))?;
677 }
678 Ok(())
679 }
680
681 fn dispatch_on_quote(&mut self, quote: QuoteTick) -> PyResult<()> {
682 if let Some(py_self) = self.python_instance()? {
683 Python::attach(|py| py_self.call_method1(py, "on_quote", (quote.into_py_any(py)?,)))?;
684 }
685 Ok(())
686 }
687
688 fn dispatch_on_trade(&mut self, trade: TradeTick) -> PyResult<()> {
689 if let Some(py_self) = self.python_instance()? {
690 Python::attach(|py| py_self.call_method1(py, "on_trade", (trade.into_py_any(py)?,)))?;
691 }
692 Ok(())
693 }
694
695 fn dispatch_on_bar(&mut self, bar: Bar) -> PyResult<()> {
696 if let Some(py_self) = self.python_instance()? {
697 Python::attach(|py| py_self.call_method1(py, "on_bar", (bar.into_py_any(py)?,)))?;
698 }
699 Ok(())
700 }
701
702 fn dispatch_on_book_deltas(&mut self, deltas: &OrderBookDeltas) -> PyResult<()> {
703 if let Some(py_self) = self.python_instance()? {
704 Python::attach(|py| {
705 py_self.call_method1(py, "on_book_deltas", (deltas.clone().into_py_any(py)?,))
706 })?;
707 }
708 Ok(())
709 }
710
711 fn dispatch_on_book_depth(&mut self, depth: &OrderBookDepth) -> PyResult<()> {
712 if let Some(py_self) = self.python_instance()? {
713 Python::attach(|py| {
714 py_self.call_method1(py, "on_book_depth", (depth.clone().into_py_any(py)?,))
715 })?;
716 }
717 Ok(())
718 }
719
720 fn dispatch_on_book(&mut self, book: &OrderBook) -> PyResult<()> {
721 if let Some(py_self) = self.python_instance()? {
722 Python::attach(|py| {
723 py_self.call_method1(py, "on_book", (book.clone().into_py_any(py)?,))
724 })?;
725 }
726 Ok(())
727 }
728
729 fn dispatch_on_mark_price(&mut self, mark_price: MarkPriceUpdate) -> PyResult<()> {
730 if let Some(py_self) = self.python_instance()? {
731 Python::attach(|py| {
732 py_self.call_method1(py, "on_mark_price", (mark_price.into_py_any(py)?,))
733 })?;
734 }
735 Ok(())
736 }
737
738 fn dispatch_on_index_price(&mut self, index_price: IndexPriceUpdate) -> PyResult<()> {
739 if let Some(py_self) = self.python_instance()? {
740 Python::attach(|py| {
741 py_self.call_method1(py, "on_index_price", (index_price.into_py_any(py)?,))
742 })?;
743 }
744 Ok(())
745 }
746
747 fn dispatch_on_funding_rate(&mut self, funding_rate: FundingRateUpdate) -> PyResult<()> {
748 if let Some(py_self) = self.python_instance()? {
749 Python::attach(|py| {
750 py_self.call_method1(py, "on_funding_rate", (funding_rate.into_py_any(py)?,))
751 })?;
752 }
753 Ok(())
754 }
755
756 fn dispatch_on_instrument_status(&mut self, data: InstrumentStatus) -> PyResult<()> {
757 if let Some(py_self) = self.python_instance()? {
758 Python::attach(|py| {
759 py_self.call_method1(py, "on_instrument_status", (data.into_py_any(py)?,))
760 })?;
761 }
762 Ok(())
763 }
764
765 fn dispatch_on_instrument_close(&mut self, update: InstrumentClose) -> PyResult<()> {
766 if let Some(py_self) = self.python_instance()? {
767 Python::attach(|py| {
768 py_self.call_method1(py, "on_instrument_close", (update.into_py_any(py)?,))
769 })?;
770 }
771 Ok(())
772 }
773
774 fn dispatch_on_option_greeks(&mut self, greeks: OptionGreeks) -> PyResult<()> {
775 if let Some(py_self) = self.python_instance()? {
776 Python::attach(|py| {
777 py_self.call_method1(py, "on_option_greeks", (greeks.into_py_any(py)?,))
778 })?;
779 }
780 Ok(())
781 }
782
783 fn dispatch_on_option_chain(&mut self, slice: &OptionChainSlice) -> PyResult<()> {
784 if let Some(py_self) = self.python_instance()? {
785 Python::attach(|py| {
786 py_self.call_method1(py, "on_option_chain", (slice.clone().into_py_any(py)?,))
787 })?;
788 }
789 Ok(())
790 }
791
792 fn dispatch_on_historical_data(&mut self, data: Py<PyAny>) -> PyResult<()> {
793 if let Some(py_self) = self.python_instance()? {
794 Python::attach(|py| py_self.call_method1(py, "on_historical_data", (data,)))?;
795 }
796 Ok(())
797 }
798
799 fn dispatch_on_historical_book_deltas(&mut self, deltas: Vec<OrderBookDelta>) -> PyResult<()> {
800 if let Some(py_self) = self.python_instance()? {
801 Python::attach(|py| {
802 let py_deltas = deltas
803 .into_iter()
804 .map(|delta| delta.into_py_any(py))
805 .collect::<PyResult<Vec<_>>>()?;
806 py_self.call_method1(py, "on_historical_book_deltas", (py_deltas,))
807 })?;
808 }
809 Ok(())
810 }
811
812 fn dispatch_on_historical_book_depth(&mut self, depths: Vec<OrderBookDepth>) -> PyResult<()> {
813 if let Some(py_self) = self.python_instance()? {
814 Python::attach(|py| {
815 let py_depths = depths
816 .into_iter()
817 .map(|depth| depth.into_py_any(py))
818 .collect::<PyResult<Vec<_>>>()?;
819 py_self.call_method1(py, "on_historical_book_depth", (py_depths,))
820 })?;
821 }
822 Ok(())
823 }
824
825 fn dispatch_on_historical_quotes(&mut self, quotes: Vec<QuoteTick>) -> PyResult<()> {
826 if let Some(py_self) = self.python_instance()? {
827 Python::attach(|py| {
828 let py_quotes = quotes
829 .into_iter()
830 .map(|quote| quote.into_py_any(py))
831 .collect::<PyResult<Vec<_>>>()?;
832 py_self.call_method1(py, "on_historical_quotes", (py_quotes,))
833 })?;
834 }
835 Ok(())
836 }
837
838 fn dispatch_on_historical_trades(&mut self, trades: Vec<TradeTick>) -> PyResult<()> {
839 if let Some(py_self) = self.python_instance()? {
840 Python::attach(|py| {
841 let py_trades = trades
842 .into_iter()
843 .map(|trade| trade.into_py_any(py))
844 .collect::<PyResult<Vec<_>>>()?;
845 py_self.call_method1(py, "on_historical_trades", (py_trades,))
846 })?;
847 }
848 Ok(())
849 }
850
851 fn dispatch_on_historical_funding_rates(
852 &mut self,
853 funding_rates: Vec<FundingRateUpdate>,
854 ) -> PyResult<()> {
855 if let Some(py_self) = self.python_instance()? {
856 Python::attach(|py| {
857 let py_funding_rates = funding_rates
858 .into_iter()
859 .map(|rate| rate.into_py_any(py))
860 .collect::<PyResult<Vec<_>>>()?;
861 py_self.call_method1(py, "on_historical_funding_rates", (py_funding_rates,))
862 })?;
863 }
864 Ok(())
865 }
866
867 fn dispatch_on_historical_bars(&mut self, bars: Vec<Bar>) -> PyResult<()> {
868 if let Some(py_self) = self.python_instance()? {
869 Python::attach(|py| {
870 let py_bars = bars
871 .into_iter()
872 .map(|bar| bar.into_py_any(py))
873 .collect::<PyResult<Vec<_>>>()?;
874 py_self.call_method1(py, "on_historical_bars", (py_bars,))
875 })?;
876 }
877 Ok(())
878 }
879
880 fn dispatch_on_historical_mark_prices(
881 &mut self,
882 mark_prices: Vec<MarkPriceUpdate>,
883 ) -> PyResult<()> {
884 if let Some(py_self) = self.python_instance()? {
885 Python::attach(|py| {
886 let py_mark_prices = mark_prices
887 .into_iter()
888 .map(|price| price.into_py_any(py))
889 .collect::<PyResult<Vec<_>>>()?;
890 py_self.call_method1(py, "on_historical_mark_prices", (py_mark_prices,))
891 })?;
892 }
893 Ok(())
894 }
895
896 fn dispatch_on_historical_index_prices(
897 &mut self,
898 index_prices: Vec<IndexPriceUpdate>,
899 ) -> PyResult<()> {
900 if let Some(py_self) = self.python_instance()? {
901 Python::attach(|py| {
902 let py_index_prices = index_prices
903 .into_iter()
904 .map(|price| price.into_py_any(py))
905 .collect::<PyResult<Vec<_>>>()?;
906 py_self.call_method1(py, "on_historical_index_prices", (py_index_prices,))
907 })?;
908 }
909 Ok(())
910 }
911
912 fn python_instance(&self) -> PyResult<Option<Py<PyAny>>> {
915 upgrade_py_weakref(
916 self.py_self.as_ref(),
917 &DataActorNative::core(&self.core).actor_id,
918 )
919 }
920}
921
922impl Deref for PyStrategyInner {
923 type Target = DataActorCore;
924
925 fn deref(&self) -> &Self::Target {
926 DataActorNative::core(&self.core)
927 }
928}
929
930impl DerefMut for PyStrategyInner {
931 fn deref_mut(&mut self) -> &mut Self::Target {
932 DataActorNative::core_mut(&mut self.core)
933 }
934}
935
936impl DataActorNative for PyStrategyInner {
937 fn core(&self) -> &DataActorCore {
938 DataActorNative::core(&self.core)
939 }
940
941 fn core_mut(&mut self) -> &mut DataActorCore {
942 DataActorNative::core_mut(&mut self.core)
943 }
944}
945
946impl StrategyNative for PyStrategyInner {
947 fn strategy_core(&self) -> &StrategyCore {
948 &self.core
949 }
950
951 fn strategy_core_mut(&mut self) -> &mut StrategyCore {
952 &mut self.core
953 }
954}
955
956impl Strategy for PyStrategyInner {
957 fn external_order_instrument_ids(&self) -> Option<Vec<InstrumentId>> {
958 self.core.config.external_order_instrument_ids.clone()
959 }
960
961 fn on_market_exit(&mut self) {
962 self.logger
963 .log_callback_error("on_market_exit", self.dispatch_on_market_exit());
964 }
965
966 fn post_market_exit(&mut self) {
967 self.logger
968 .log_callback_error("post_market_exit", self.dispatch_post_market_exit());
969 }
970
971 fn on_order_initialized(&mut self, event: OrderInitialized) {
972 self.logger.log_callback_error(
973 "on_order_initialized",
974 self.dispatch_on_order_initialized(event),
975 );
976 }
977
978 fn on_order_event(&mut self, event: OrderEventAny) {
979 self.logger
980 .log_callback_error("on_order_event", self.dispatch_on_order_event(event));
981 }
982
983 fn on_order_denied(&mut self, event: OrderDenied) {
984 self.logger
985 .log_callback_error("on_order_denied", self.dispatch_on_order_denied(event));
986 }
987
988 fn on_order_emulated(&mut self, event: OrderEmulated) {
989 self.logger
990 .log_callback_error("on_order_emulated", self.dispatch_on_order_emulated(event));
991 }
992
993 fn on_order_released(&mut self, event: OrderReleased) {
994 self.logger
995 .log_callback_error("on_order_released", self.dispatch_on_order_released(event));
996 }
997
998 fn on_order_submitted(&mut self, event: OrderSubmitted) {
999 self.logger.log_callback_error(
1000 "on_order_submitted",
1001 self.dispatch_on_order_submitted(event),
1002 );
1003 }
1004
1005 fn on_order_rejected(&mut self, event: OrderRejected) {
1006 self.logger
1007 .log_callback_error("on_order_rejected", self.dispatch_on_order_rejected(event));
1008 }
1009
1010 fn on_order_accepted(&mut self, event: OrderAccepted) {
1011 self.logger
1012 .log_callback_error("on_order_accepted", self.dispatch_on_order_accepted(event));
1013 }
1014
1015 fn on_order_expired(&mut self, event: OrderExpired) {
1016 self.logger
1017 .log_callback_error("on_order_expired", self.dispatch_on_order_expired(event));
1018 }
1019
1020 fn on_order_triggered(&mut self, event: OrderTriggered) {
1021 self.logger.log_callback_error(
1022 "on_order_triggered",
1023 self.dispatch_on_order_triggered(event),
1024 );
1025 }
1026
1027 fn on_order_pending_update(&mut self, event: OrderPendingUpdate) {
1028 self.logger.log_callback_error(
1029 "on_order_pending_update",
1030 self.dispatch_on_order_pending_update(event),
1031 );
1032 }
1033
1034 fn on_order_pending_cancel(&mut self, event: OrderPendingCancel) {
1035 self.logger.log_callback_error(
1036 "on_order_pending_cancel",
1037 self.dispatch_on_order_pending_cancel(event),
1038 );
1039 }
1040
1041 fn on_order_modify_rejected(&mut self, event: OrderModifyRejected) {
1042 self.logger.log_callback_error(
1043 "on_order_modify_rejected",
1044 self.dispatch_on_order_modify_rejected(event),
1045 );
1046 }
1047
1048 fn on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {
1049 self.logger.log_callback_error(
1050 "on_order_cancel_rejected",
1051 self.dispatch_on_order_cancel_rejected(event),
1052 );
1053 }
1054
1055 fn on_order_updated(&mut self, event: OrderUpdated) {
1056 self.logger
1057 .log_callback_error("on_order_updated", self.dispatch_on_order_updated(&event));
1058 }
1059
1060 fn on_order_canceled(&mut self, event: &OrderCanceled) {
1061 self.logger
1062 .log_callback_error("on_order_canceled", self.dispatch_on_order_canceled(*event));
1063 }
1064
1065 fn on_order_filled(&mut self, event: &OrderFilled) {
1066 self.logger
1067 .log_callback_error("on_order_filled", self.dispatch_on_order_filled(event));
1068 }
1069
1070 fn on_order_fill_voided(&mut self, event: &OrderFillVoided) {
1071 self.logger.log_callback_error(
1072 "on_order_fill_voided",
1073 self.dispatch_on_order_fill_voided(event),
1074 );
1075 }
1076
1077 fn on_position_opened(&mut self, event: PositionOpened) {
1078 self.logger.log_callback_error(
1079 "on_position_opened",
1080 self.dispatch_on_position_opened(event),
1081 );
1082 }
1083
1084 fn on_position_event(&mut self, event: PositionEvent) {
1085 self.logger
1086 .log_callback_error("on_position_event", self.dispatch_on_position_event(event));
1087 }
1088
1089 fn on_position_changed(&mut self, event: PositionChanged) {
1090 self.logger.log_callback_error(
1091 "on_position_changed",
1092 self.dispatch_on_position_changed(event),
1093 );
1094 }
1095
1096 fn on_position_closed(&mut self, event: PositionClosed) {
1097 self.logger.log_callback_error(
1098 "on_position_closed",
1099 self.dispatch_on_position_closed(event),
1100 );
1101 }
1102}
1103
1104impl DataActor for PyStrategyInner {
1105 fn on_start(&mut self) -> anyhow::Result<()> {
1106 Strategy::on_start(self)?;
1107 self.dispatch_on_start()
1108 .map_err(|e| anyhow::anyhow!("Python on_start failed:\n{}", format_exception(&e)))
1109 }
1110
1111 fn on_stop(&mut self) -> anyhow::Result<()> {
1112 self.dispatch_on_stop()
1113 .map_err(|e| anyhow::anyhow!("Python on_stop failed:\n{}", format_exception(&e)))
1114 }
1115
1116 fn on_resume(&mut self) -> anyhow::Result<()> {
1117 self.dispatch_on_resume()
1118 .map_err(|e| anyhow::anyhow!("Python on_resume failed:\n{}", format_exception(&e)))
1119 }
1120
1121 fn on_reset(&mut self) -> anyhow::Result<()> {
1122 self.dispatch_on_reset()
1123 .map_err(|e| anyhow::anyhow!("Python on_reset failed:\n{}", format_exception(&e)))
1124 }
1125
1126 fn on_dispose(&mut self) -> anyhow::Result<()> {
1127 self.dispatch_on_dispose()
1128 .map_err(|e| anyhow::anyhow!("Python on_dispose failed:\n{}", format_exception(&e)))
1129 }
1130
1131 fn on_degrade(&mut self) -> anyhow::Result<()> {
1132 self.dispatch_on_degrade()
1133 .map_err(|e| anyhow::anyhow!("Python on_degrade failed:\n{}", format_exception(&e)))
1134 }
1135
1136 fn on_fault(&mut self) -> anyhow::Result<()> {
1137 self.dispatch_on_fault()
1138 .map_err(|e| anyhow::anyhow!("Python on_fault failed:\n{}", format_exception(&e)))
1139 }
1140
1141 fn on_save(&self) -> anyhow::Result<IndexMap<String, Vec<u8>>> {
1142 self.dispatch_on_save()
1143 .map_err(|e| anyhow::anyhow!("Python on_save failed:\n{}", format_exception(&e)))
1144 }
1145
1146 fn on_load(&mut self, state: IndexMap<String, Vec<u8>>) -> anyhow::Result<()> {
1147 self.dispatch_on_load(&state)
1148 .map_err(|e| anyhow::anyhow!("Python on_load failed:\n{}", format_exception(&e)))
1149 }
1150
1151 fn on_time_event(&mut self, event: &TimeEvent) -> anyhow::Result<()> {
1152 route_time_event(self, event);
1153 self.dispatch_on_time_event(event)
1154 .map_err(|e| anyhow::anyhow!("Python on_time_event failed:\n{}", format_exception(&e)))
1155 }
1156
1157 #[allow(unused_variables)]
1158 fn on_data(&mut self, data: &CustomData) -> anyhow::Result<()> {
1159 Python::attach(|py| {
1160 let py_data: Py<PyAny> = Py::new(py, data.clone())?.into_any();
1161 self.dispatch_on_data(py_data)
1162 .map_err(|e| anyhow::anyhow!("Python on_data failed:\n{}", format_exception(&e)))
1163 })
1164 }
1165
1166 fn on_signal(&mut self, signal: &Signal) -> anyhow::Result<()> {
1167 self.dispatch_on_signal(signal)
1168 .map_err(|e| anyhow::anyhow!("Python on_signal failed:\n{}", format_exception(&e)))
1169 }
1170
1171 fn on_queue_state(&mut self, event: &QueueStateChanged) -> anyhow::Result<()> {
1172 self.dispatch_on_queue_state(event)
1173 .map_err(|e| anyhow::anyhow!("Python on_queue_state failed:\n{}", format_exception(&e)))
1174 }
1175
1176 fn on_socket_state(&mut self, event: &SocketStateChanged) -> anyhow::Result<()> {
1177 self.dispatch_on_socket_state(event).map_err(|e| {
1178 anyhow::anyhow!("Python on_socket_state failed:\n{}", format_exception(&e))
1179 })
1180 }
1181
1182 fn on_instrument(&mut self, instrument: &InstrumentAny) -> anyhow::Result<()> {
1183 Python::attach(|py| {
1184 let py_instrument = instrument_any_to_pyobject(py, instrument.clone())
1185 .map_err(|e| anyhow::anyhow!("Failed to convert InstrumentAny to Python: {e}"))?;
1186 self.dispatch_on_instrument(py_instrument).map_err(|e| {
1187 anyhow::anyhow!("Python on_instrument failed:\n{}", format_exception(&e))
1188 })
1189 })
1190 }
1191
1192 fn on_quote(&mut self, quote: &QuoteTick) -> anyhow::Result<()> {
1193 self.dispatch_on_quote(*quote)
1194 .map_err(|e| anyhow::anyhow!("Python on_quote failed:\n{}", format_exception(&e)))
1195 }
1196
1197 fn on_trade(&mut self, tick: &TradeTick) -> anyhow::Result<()> {
1198 self.dispatch_on_trade(*tick)
1199 .map_err(|e| anyhow::anyhow!("Python on_trade failed:\n{}", format_exception(&e)))
1200 }
1201
1202 fn on_bar(&mut self, bar: &Bar) -> anyhow::Result<()> {
1203 self.dispatch_on_bar(*bar)
1204 .map_err(|e| anyhow::anyhow!("Python on_bar failed:\n{}", format_exception(&e)))
1205 }
1206
1207 fn on_book_deltas(&mut self, deltas: &OrderBookDeltas) -> anyhow::Result<()> {
1208 self.dispatch_on_book_deltas(deltas)
1209 .map_err(|e| anyhow::anyhow!("Python on_book_deltas failed:\n{}", format_exception(&e)))
1210 }
1211
1212 fn on_book_depth(&mut self, depth: &OrderBookDepth) -> anyhow::Result<()> {
1213 self.dispatch_on_book_depth(depth)
1214 .map_err(|e| anyhow::anyhow!("Python on_book_depth failed:\n{}", format_exception(&e)))
1215 }
1216
1217 fn on_book(&mut self, order_book: &OrderBook) -> anyhow::Result<()> {
1218 self.dispatch_on_book(order_book)
1219 .map_err(|e| anyhow::anyhow!("Python on_book failed:\n{}", format_exception(&e)))
1220 }
1221
1222 fn on_mark_price(&mut self, mark_price: &MarkPriceUpdate) -> anyhow::Result<()> {
1223 self.dispatch_on_mark_price(*mark_price)
1224 .map_err(|e| anyhow::anyhow!("Python on_mark_price failed:\n{}", format_exception(&e)))
1225 }
1226
1227 fn on_index_price(&mut self, index_price: &IndexPriceUpdate) -> anyhow::Result<()> {
1228 self.dispatch_on_index_price(*index_price)
1229 .map_err(|e| anyhow::anyhow!("Python on_index_price failed:\n{}", format_exception(&e)))
1230 }
1231
1232 fn on_funding_rate(&mut self, funding_rate: &FundingRateUpdate) -> anyhow::Result<()> {
1233 self.dispatch_on_funding_rate(*funding_rate).map_err(|e| {
1234 anyhow::anyhow!("Python on_funding_rate failed:\n{}", format_exception(&e))
1235 })
1236 }
1237
1238 fn on_instrument_status(&mut self, data: &InstrumentStatus) -> anyhow::Result<()> {
1239 self.dispatch_on_instrument_status(*data).map_err(|e| {
1240 anyhow::anyhow!(
1241 "Python on_instrument_status failed:\n{}",
1242 format_exception(&e)
1243 )
1244 })
1245 }
1246
1247 fn on_instrument_close(&mut self, update: &InstrumentClose) -> anyhow::Result<()> {
1248 self.dispatch_on_instrument_close(*update).map_err(|e| {
1249 anyhow::anyhow!(
1250 "Python on_instrument_close failed:\n{}",
1251 format_exception(&e)
1252 )
1253 })
1254 }
1255
1256 fn on_option_greeks(&mut self, greeks: &OptionGreeks) -> anyhow::Result<()> {
1257 self.dispatch_on_option_greeks(*greeks).map_err(|e| {
1258 anyhow::anyhow!("Python on_option_greeks failed:\n{}", format_exception(&e))
1259 })
1260 }
1261
1262 fn on_option_chain(&mut self, slice: &OptionChainSlice) -> anyhow::Result<()> {
1263 self.dispatch_on_option_chain(slice).map_err(|e| {
1264 anyhow::anyhow!("Python on_option_chain failed:\n{}", format_exception(&e))
1265 })
1266 }
1267
1268 fn on_historical_data(&mut self, data: &dyn Any) -> anyhow::Result<()> {
1269 Python::attach(|py| {
1270 let py_data: Py<PyAny> = if let Some(custom_data) = data.downcast_ref::<CustomData>() {
1271 Py::new(py, custom_data.clone())?.into_any()
1272 } else if let Some(custom_data) = data.downcast_ref::<Vec<CustomData>>() {
1273 custom_data.clone().into_py_any(py)?
1274 } else {
1275 anyhow::bail!("Failed to convert historical data to Python: unsupported type");
1276 };
1277
1278 self.dispatch_on_historical_data(py_data).map_err(|e| {
1279 anyhow::anyhow!(
1280 "Python on_historical_data failed:\n{}",
1281 format_exception(&e)
1282 )
1283 })
1284 })
1285 }
1286
1287 fn on_historical_book_deltas(&mut self, deltas: &[OrderBookDelta]) -> anyhow::Result<()> {
1288 self.dispatch_on_historical_book_deltas(deltas.to_vec())
1289 .map_err(|e| {
1290 anyhow::anyhow!(
1291 "Python on_historical_book_deltas failed:\n{}",
1292 format_exception(&e)
1293 )
1294 })
1295 }
1296
1297 fn on_historical_book_depth(&mut self, depths: &[OrderBookDepth]) -> anyhow::Result<()> {
1298 self.dispatch_on_historical_book_depth(depths.to_vec())
1299 .map_err(|e| {
1300 anyhow::anyhow!(
1301 "Python on_historical_book_depth failed:\n{}",
1302 format_exception(&e)
1303 )
1304 })
1305 }
1306
1307 fn on_historical_quotes(&mut self, quotes: &[QuoteTick]) -> anyhow::Result<()> {
1308 self.dispatch_on_historical_quotes(quotes.to_vec())
1309 .map_err(|e| {
1310 anyhow::anyhow!(
1311 "Python on_historical_quotes failed:\n{}",
1312 format_exception(&e)
1313 )
1314 })
1315 }
1316
1317 fn on_historical_trades(&mut self, trades: &[TradeTick]) -> anyhow::Result<()> {
1318 self.dispatch_on_historical_trades(trades.to_vec())
1319 .map_err(|e| {
1320 anyhow::anyhow!(
1321 "Python on_historical_trades failed:\n{}",
1322 format_exception(&e)
1323 )
1324 })
1325 }
1326
1327 fn on_historical_funding_rates(
1328 &mut self,
1329 funding_rates: &[FundingRateUpdate],
1330 ) -> anyhow::Result<()> {
1331 self.dispatch_on_historical_funding_rates(funding_rates.to_vec())
1332 .map_err(|e| {
1333 anyhow::anyhow!(
1334 "Python on_historical_funding_rates failed:\n{}",
1335 format_exception(&e)
1336 )
1337 })
1338 }
1339
1340 fn on_historical_bars(&mut self, bars: &[Bar]) -> anyhow::Result<()> {
1341 self.dispatch_on_historical_bars(bars.to_vec())
1342 .map_err(|e| {
1343 anyhow::anyhow!(
1344 "Python on_historical_bars failed:\n{}",
1345 format_exception(&e)
1346 )
1347 })
1348 }
1349
1350 fn on_historical_mark_prices(&mut self, mark_prices: &[MarkPriceUpdate]) -> anyhow::Result<()> {
1351 self.dispatch_on_historical_mark_prices(mark_prices.to_vec())
1352 .map_err(|e| {
1353 anyhow::anyhow!(
1354 "Python on_historical_mark_prices failed:\n{}",
1355 format_exception(&e)
1356 )
1357 })
1358 }
1359
1360 fn on_historical_index_prices(
1361 &mut self,
1362 index_prices: &[IndexPriceUpdate],
1363 ) -> anyhow::Result<()> {
1364 self.dispatch_on_historical_index_prices(index_prices.to_vec())
1365 .map_err(|e| {
1366 anyhow::anyhow!(
1367 "Python on_historical_index_prices failed:\n{}",
1368 format_exception(&e)
1369 )
1370 })
1371 }
1372}
1373
1374fn state_to_pydict(py: Python<'_>, state: &IndexMap<String, Vec<u8>>) -> PyResult<Py<PyDict>> {
1375 let py_state = PyDict::new(py);
1376 for (key, value) in state {
1377 py_state.set_item(key, PyBytes::new(py, value))?;
1378 }
1379 Ok(py_state.unbind())
1380}
1381
1382fn pydict_to_state(state: &Bound<'_, PyDict>) -> PyResult<IndexMap<String, Vec<u8>>> {
1383 let mut rust_state = IndexMap::with_capacity(state.len());
1384 for (key, value) in state.iter() {
1385 rust_state.insert(key.extract()?, value.extract()?);
1386 }
1387 Ok(rust_state)
1388}
1389
1390#[allow(non_camel_case_types)]
1392#[pyo3::pyclass(
1393 module = "nautilus_trader.trading",
1394 name = "Strategy",
1395 unsendable,
1396 subclass,
1397 weakref
1398)]
1399#[pyo3_stub_gen::derive::gen_stub_pyclass(module = "nautilus_trader.trading")]
1400pub struct PyStrategy {
1401 inner: Rc<UnsafeCell<PyStrategyInner>>,
1402}
1403
1404impl Debug for PyStrategy {
1405 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1406 f.debug_struct(stringify!(PyStrategy))
1407 .field("inner", &self.inner())
1408 .finish()
1409 }
1410}
1411
1412impl PyStrategy {
1413 #[inline]
1414 #[allow(unsafe_code)]
1415 pub(crate) fn inner(&self) -> &PyStrategyInner {
1416 unsafe { &*self.inner.get() }
1419 }
1420
1421 #[inline]
1422 #[allow(unsafe_code, clippy::mut_from_ref)]
1423 pub(crate) fn inner_mut(&self) -> &mut PyStrategyInner {
1424 unsafe { &mut *self.inner.get() }
1427 }
1428}
1429
1430impl PyStrategy {
1431 pub fn new_checked(config: Option<StrategyConfig>) -> CorrectnessResult<Self> {
1438 let config = config.unwrap_or_default();
1439 let core = StrategyCore::new_checked(config)?;
1440 let clock = PyClock::new_test();
1441 let logger = PyLogger::new(core.actor.actor_id.as_str());
1442
1443 let inner = PyStrategyInner {
1444 core,
1445 py_self: None,
1446 config: None,
1447 clock,
1448 logger,
1449 };
1450
1451 Ok(Self {
1452 inner: Rc::new(UnsafeCell::new(inner)),
1453 })
1454 }
1455
1456 #[must_use]
1463 pub fn new(config: Option<StrategyConfig>) -> Self {
1464 Self::new_checked(config).expect_display(FAILED)
1465 }
1466
1467 pub fn set_python_instance(&mut self, py_obj: &Bound<'_, PyAny>) -> PyResult<()> {
1476 self.inner_mut().py_self = Some(PyWeakrefReference::new(py_obj)?.unbind());
1477 Ok(())
1478 }
1479
1480 pub fn set_config(&mut self, config: Option<Py<PyAny>>) {
1485 self.inner_mut().config = config;
1486 }
1487
1488 pub fn set_external_order_instrument_ids(
1490 &mut self,
1491 external_order_instrument_ids: Option<Vec<InstrumentId>>,
1492 ) {
1493 self.inner_mut().core.config.external_order_instrument_ids = external_order_instrument_ids;
1494 }
1495
1496 #[must_use]
1498 pub fn external_order_instrument_ids(&self) -> Option<Vec<InstrumentId>> {
1499 self.inner().external_order_instrument_ids()
1500 }
1501
1502 pub fn set_actor_id(&mut self, actor_id: ActorId) {
1506 let inner = self.inner_mut();
1507 inner.core.actor.config.actor_id = Some(actor_id);
1508 inner.core.actor.actor_id = actor_id;
1509 inner.logger = PyLogger::new(actor_id.as_str());
1510 }
1511
1512 pub fn set_strategy_id(&mut self, strategy_id: StrategyId) -> anyhow::Result<()> {
1521 let inner = self.inner_mut();
1522 inner.core.change_id(strategy_id)?;
1523 inner.logger = PyLogger::new(inner.core.actor.actor_id.as_str());
1524 Ok(())
1525 }
1526
1527 pub fn set_order_id_tag(&mut self, order_id_tag: &str) -> anyhow::Result<()> {
1534 let inner = self.inner_mut();
1535 inner.core.change_order_id_tag(order_id_tag)?;
1536 inner.logger = PyLogger::new(inner.core.actor.actor_id.as_str());
1537 Ok(())
1538 }
1539
1540 pub fn set_log_events(&mut self, log_events: bool) {
1542 let inner = self.inner_mut();
1543 inner.core.actor.config.log_events = log_events;
1544 }
1545
1546 pub fn set_log_commands(&mut self, log_commands: bool) {
1548 let inner = self.inner_mut();
1549 inner.core.actor.config.log_commands = log_commands;
1550 }
1551
1552 #[must_use]
1557 pub fn strategy_id(&self) -> StrategyId {
1558 StrategyId::new(self.inner().core.actor.actor_id.inner())
1559 }
1560
1561 #[must_use]
1563 pub fn configured_strategy_id(&self) -> Option<StrategyId> {
1564 self.inner().core.strategy_id()
1565 }
1566
1567 #[must_use]
1569 pub fn order_id_tag(&self) -> Option<String> {
1570 self.inner().core.order_id_tag().map(str::to_string)
1571 }
1572
1573 #[must_use]
1575 pub fn is_registered(&self) -> bool {
1576 self.inner().core.actor.is_registered()
1577 }
1578
1579 pub fn register(
1585 &mut self,
1586 trader_id: TraderId,
1587 clock: Rc<RefCell<dyn Clock>>,
1588 cache: Rc<RefCell<Cache>>,
1589 portfolio: Rc<RefCell<Portfolio>>,
1590 ) -> anyhow::Result<()> {
1591 let inner = self.inner_mut();
1592 inner.core.register(trader_id, clock, cache, portfolio)?;
1593
1594 inner.clock = PyClock::from_rc(inner.core.actor.clock_rc());
1595
1596 let actor_id = inner.core.actor.actor_id.inner();
1597 let callback = TimeEventCallback::from(move |event: TimeEvent| {
1598 if let Some(mut strategy) = try_get_actor_unchecked::<PyStrategyInner>(&actor_id) {
1599 strategy.handle_time_event(&event);
1600 } else {
1601 log::error!("Strategy {actor_id} not found for time event handling");
1602 }
1603 });
1604
1605 inner.clock.inner_mut().register_default_handler(callback);
1606
1607 Component::initialize(inner)
1608 }
1609
1610 pub fn register_in_global_registries(&self) -> PyResult<()> {
1620 let inner = self.inner();
1621 let component_id = Component::component_id(inner);
1622 let actor_id = Actor::id(inner);
1623
1624 let Some(wrapper) = inner.python_instance()? else {
1625 return Err(to_pyruntime_err(format!(
1626 "Cannot register strategy {actor_id} without a Python wrapper, call `set_python_instance` first"
1627 )));
1628 };
1629
1630 let inner_ref: Rc<UnsafeCell<PyStrategyInner>> = self.inner.clone();
1631
1632 let component_trait_ref: Rc<UnsafeCell<dyn Component>> = inner_ref.clone();
1633 with_component_registry(|registry| {
1634 registry.insert(component_id.inner(), component_trait_ref);
1635 });
1636
1637 let actor_trait_ref: Rc<UnsafeCell<dyn Actor>> = inner_ref;
1638 with_actor_registry(|registry| registry.insert(actor_id, actor_trait_ref));
1639
1640 retain_python_wrapper(component_id, wrapper, inner.core.actor.message_bus());
1641
1642 Ok(())
1643 }
1644}
1645
1646#[pyo3::pymethods]
1647#[pyo3_stub_gen::derive::gen_stub_pymethods]
1648#[allow(
1649 clippy::large_types_passed_by_value,
1650 reason = "PyO3 callbacks accept Python-owned event values"
1651)]
1652#[expect(
1653 clippy::unused_self,
1654 reason = "default PyO3 callbacks must remain instance methods"
1655)]
1656impl PyStrategy {
1657 #[new]
1669 #[pyo3(signature = (config=None))]
1670 fn py_new(config: Option<Py<PyAny>>) -> PyResult<Self> {
1671 let strategy_config = config
1672 .as_ref()
1673 .and_then(|obj| Python::attach(|py| obj.extract::<StrategyConfig>(py).ok()));
1674 let mut strategy = Self::new_checked(strategy_config).map_err(to_pyvalue_err)?;
1675 strategy.set_config(config);
1676 Ok(strategy)
1677 }
1678
1679 #[pyo3(signature = (config=None))]
1681 fn __init__(slf: &Bound<'_, Self>, config: Option<Py<PyAny>>) -> PyResult<()> {
1682 {
1683 let mut borrowed = slf.borrow_mut();
1684 borrowed.set_python_instance(slf.as_any())?;
1685 if config.is_some() {
1687 borrowed.set_config(config);
1688 }
1689 }
1690
1691 if !has_configured_strategy_id(slf) {
1692 let py_type = slf.get_type();
1693 let type_name = py_type.name()?;
1694 let actor_id = class_derived_actor_id(slf, type_name.to_str()?)?;
1695 slf.borrow_mut().set_actor_id(actor_id);
1696 }
1697
1698 Ok(())
1699 }
1700
1701 #[getter]
1702 #[pyo3(name = "trader_id")]
1703 fn py_trader_id(&self) -> Option<TraderId> {
1704 self.inner().core.trader_id()
1705 }
1706
1707 #[getter]
1708 #[pyo3(name = "strategy_id")]
1709 fn py_strategy_id(&self) -> StrategyId {
1710 self.strategy_id()
1711 }
1712
1713 #[getter]
1714 #[pyo3(name = "config")]
1715 fn py_config(&self, py: Python<'_>) -> Option<Py<PyAny>> {
1716 self.inner()
1717 .config
1718 .as_ref()
1719 .map(|config| config.clone_ref(py))
1720 }
1721
1722 #[getter]
1723 #[pyo3(name = "clock")]
1724 fn py_clock(&self) -> PyResult<PyClock> {
1725 let inner = self.inner();
1726 if inner.core.actor.is_registered() {
1727 Ok(inner.clock.clone())
1728 } else {
1729 Err(to_pyruntime_err(
1730 "Strategy must be registered with a trader before accessing clock",
1731 ))
1732 }
1733 }
1734
1735 #[getter]
1736 #[pyo3(name = "cache")]
1737 fn py_cache(&self) -> PyResult<PyCache> {
1738 let inner = self.inner();
1739 if inner.core.actor.is_registered() {
1740 Ok(PyCache::from_rc(inner.core.actor.cache_rc()))
1741 } else {
1742 Err(to_pyruntime_err(
1743 "Strategy must be registered with a trader before accessing cache",
1744 ))
1745 }
1746 }
1747
1748 #[pyo3(name = "set_external_order_instrument_ids")]
1758 fn py_set_external_order_instrument_ids(
1759 &mut self,
1760 instrument_ids: Vec<InstrumentId>,
1761 ) -> PyResult<()> {
1762 Strategy::set_external_order_instrument_ids(self.inner_mut(), instrument_ids)
1763 .map_err(to_pyruntime_err)
1764 }
1765
1766 #[getter]
1767 #[pyo3(name = "portfolio")]
1768 fn py_portfolio(&self) -> PyResult<PyPortfolio> {
1769 let inner = self.inner();
1770 if inner.core.actor.is_registered() {
1771 Ok(PyPortfolio::from_rc(inner.portfolio_rc()))
1772 } else {
1773 Err(to_pyruntime_err(
1774 "Strategy must be registered with a trader before accessing portfolio",
1775 ))
1776 }
1777 }
1778
1779 #[getter]
1780 #[pyo3(name = "order_factory")]
1781 fn py_order_factory(&self) -> PyResult<PyOrderFactory> {
1782 let inner = self.inner();
1783 if inner.core.actor.is_registered() {
1784 Ok(PyOrderFactory::from_rc(inner.order_factory_rc()))
1785 } else {
1786 Err(to_pyruntime_err(
1787 "Strategy must be registered with a trader before accessing order_factory",
1788 ))
1789 }
1790 }
1791
1792 #[getter]
1793 #[pyo3(name = "log")]
1794 fn py_log(&self) -> PyLogger {
1795 self.inner().logger.clone()
1796 }
1797
1798 #[pyo3(name = "state")]
1799 fn py_state(&self) -> ComponentState {
1800 self.inner().core.actor.state()
1801 }
1802
1803 #[pyo3(name = "is_ready")]
1804 fn py_is_ready(&self) -> bool {
1805 Component::is_ready(self.inner())
1806 }
1807
1808 #[pyo3(name = "is_running")]
1809 fn py_is_running(&self) -> bool {
1810 Component::is_running(self.inner())
1811 }
1812
1813 #[pyo3(name = "is_stopped")]
1814 fn py_is_stopped(&self) -> bool {
1815 Component::is_stopped(self.inner())
1816 }
1817
1818 #[pyo3(name = "is_disposed")]
1819 fn py_is_disposed(&self) -> bool {
1820 Component::is_disposed(self.inner())
1821 }
1822
1823 #[pyo3(name = "is_degraded")]
1824 fn py_is_degraded(&self) -> bool {
1825 Component::is_degraded(self.inner())
1826 }
1827
1828 #[pyo3(name = "is_faulted")]
1829 fn py_is_faulted(&self) -> bool {
1830 Component::is_faulted(self.inner())
1831 }
1832
1833 #[pyo3(name = "start")]
1834 fn py_start(&mut self) -> PyResult<()> {
1835 Component::start(self.inner_mut()).map_err(to_pyruntime_err)
1836 }
1837
1838 #[pyo3(name = "stop")]
1839 fn py_stop(&mut self) -> PyResult<()> {
1840 let inner = self.inner_mut();
1841 if Strategy::stop(inner) {
1842 Component::stop(inner).map_err(to_pyruntime_err)
1843 } else {
1844 Ok(())
1845 }
1846 }
1847
1848 #[pyo3(name = "market_exit")]
1849 fn py_market_exit(&mut self) -> PyResult<()> {
1850 Strategy::market_exit(self.inner_mut()).map_err(to_pyruntime_err)
1851 }
1852
1853 #[pyo3(name = "is_exiting")]
1854 fn py_is_exiting(&self) -> bool {
1855 Strategy::is_exiting(self.inner())
1856 }
1857
1858 #[pyo3(name = "save")]
1859 fn py_save(&self, py: Python<'_>) -> PyResult<Py<PyDict>> {
1860 let state = DataActor::on_save(self.inner()).map_err(to_pyruntime_err)?;
1861 state_to_pydict(py, &state)
1862 }
1863
1864 #[pyo3(name = "load")]
1865 fn py_load(&mut self, state: &Bound<'_, PyDict>) -> PyResult<()> {
1866 let state = pydict_to_state(state)?;
1867 DataActor::on_load(self.inner_mut(), state).map_err(to_pyruntime_err)
1868 }
1869
1870 #[pyo3(name = "publish_data")]
1871 fn py_publish_data(&self, data_type: &DataType, data: &CustomData) -> PyResult<()> {
1872 self.ensure_registered_for_data()?;
1873 DataActor::publish_data(self.inner(), data_type, data);
1874 Ok(())
1875 }
1876
1877 #[pyo3(name = "publish_signal")]
1878 #[pyo3(signature = (name, value, ts_event=0))]
1879 #[expect(
1880 clippy::needless_pass_by_value,
1881 reason = "PyO3 accepts an owned PyAny handle for Python signal values"
1882 )]
1883 fn py_publish_signal(
1884 &self,
1885 py: Python<'_>,
1886 name: &str,
1887 value: Py<PyAny>,
1888 ts_event: u64,
1889 ) -> PyResult<()> {
1890 self.ensure_registered_for_data()?;
1891 let value_str: String = value.bind(py).str()?.extract()?;
1892 DataActor::publish_signal(self.inner(), name, value_str, UnixNanos::from(ts_event));
1893 Ok(())
1894 }
1895
1896 #[pyo3(name = "add_synthetic")]
1897 fn py_add_synthetic(&self, synthetic: SyntheticInstrument) -> PyResult<()> {
1898 self.ensure_registered_for_data()?;
1899 DataActor::add_synthetic(self.inner(), synthetic).map_err(to_pyvalue_err)
1900 }
1901
1902 #[pyo3(name = "update_synthetic")]
1903 fn py_update_synthetic(&self, synthetic: SyntheticInstrument) -> PyResult<()> {
1904 self.ensure_registered_for_data()?;
1905 DataActor::update_synthetic(self.inner(), synthetic).map_err(to_pyvalue_err)
1906 }
1907
1908 #[pyo3(name = "resume")]
1909 fn py_resume(&mut self) -> PyResult<()> {
1910 Component::resume(self.inner_mut()).map_err(to_pyruntime_err)
1911 }
1912
1913 #[pyo3(name = "reset")]
1914 fn py_reset(&mut self) -> PyResult<()> {
1915 Component::reset(self.inner_mut()).map_err(to_pyruntime_err)
1916 }
1917
1918 #[pyo3(name = "dispose")]
1919 fn py_dispose(&mut self) -> PyResult<()> {
1920 Component::dispose(self.inner_mut()).map_err(to_pyruntime_err)
1921 }
1922
1923 #[pyo3(name = "degrade")]
1924 fn py_degrade(&mut self) -> PyResult<()> {
1925 Component::degrade(self.inner_mut()).map_err(to_pyruntime_err)
1926 }
1927
1928 #[pyo3(name = "fault")]
1929 fn py_fault(&mut self) -> PyResult<()> {
1930 Component::fault(self.inner_mut()).map_err(to_pyruntime_err)
1931 }
1932
1933 #[pyo3(name = "shutdown_system")]
1934 #[pyo3(signature = (reason=None))]
1935 fn py_shutdown_system(&self, reason: Option<String>) -> PyResult<()> {
1936 let inner = self.inner();
1937 if !inner.core.actor.is_registered() {
1938 return Err(to_pyruntime_err(
1939 "Strategy must be registered with a trader before shutting down the system",
1940 ));
1941 }
1942
1943 DataActor::shutdown_system(inner, reason);
1944 Ok(())
1945 }
1946
1947 #[getter]
1948 #[pyo3(name = "registered_indicators")]
1949 fn py_registered_indicators(&self, py: Python<'_>) -> PyResult<Py<PyList>> {
1950 let inner = self.inner();
1951 registered_python_indicators(
1952 py,
1953 DataActorNative::core(&inner.core).registered_indicators(),
1954 )
1955 }
1956
1957 #[pyo3(name = "indicators_initialized")]
1958 fn py_indicators_initialized(&self, _py: Python<'_>) -> PyResult<bool> {
1959 let inner = self.inner();
1960 DataActorNative::core(&inner.core)
1961 .indicators_initialized()
1962 .map_err(to_pyruntime_err)
1963 }
1964
1965 #[pyo3(name = "register_indicator_for_quote_ticks")]
1966 fn py_register_indicator_for_quote_ticks(
1967 &mut self,
1968 py: Python<'_>,
1969 instrument_id: InstrumentId,
1970 indicator: Py<PyAny>,
1971 ) {
1972 let indicator = wrap_python_indicator(py, indicator);
1973 let inner = self.inner_mut();
1974 DataActorNative::core_mut(&mut inner.core)
1975 .register_indicator_for_quote_ticks(instrument_id, indicator);
1976 }
1977
1978 #[pyo3(name = "register_indicator_for_trade_ticks")]
1979 fn py_register_indicator_for_trade_ticks(
1980 &mut self,
1981 py: Python<'_>,
1982 instrument_id: InstrumentId,
1983 indicator: Py<PyAny>,
1984 ) {
1985 let indicator = wrap_python_indicator(py, indicator);
1986 let inner = self.inner_mut();
1987 DataActorNative::core_mut(&mut inner.core)
1988 .register_indicator_for_trade_ticks(instrument_id, indicator);
1989 }
1990
1991 #[pyo3(name = "register_indicator_for_bars")]
1992 fn py_register_indicator_for_bars(
1993 &mut self,
1994 py: Python<'_>,
1995 bar_type: BarType,
1996 indicator: Py<PyAny>,
1997 ) {
1998 let indicator = wrap_python_indicator(py, indicator);
1999 let inner = self.inner_mut();
2000 DataActorNative::core_mut(&mut inner.core).register_indicator_for_bars(bar_type, indicator);
2001 }
2002
2003 #[pyo3(name = "submit_order")]
2004 #[pyo3(signature = (order, position_id=None, client_id=None, params=None))]
2005 fn py_submit_order(
2006 &mut self,
2007 py: Python<'_>,
2008 order: Py<PyAny>,
2009 position_id: Option<PositionId>,
2010 client_id: Option<ClientId>,
2011 params: Option<Py<PyDict>>,
2012 ) -> PyResult<()> {
2013 let order = pyobject_to_order_any(py, order)?;
2014 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2015 match params {
2016 Some(dict) => from_pydict(py, &dict),
2017 None => Ok(None),
2018 }
2019 })?;
2020 let inner = self.inner_mut();
2021
2022 Strategy::submit_order(inner, order, position_id, client_id, params_map)
2023 .map_err(to_pyruntime_err)
2024 }
2025
2026 #[pyo3(name = "submit_order_list")]
2027 #[pyo3(signature = (order_list, position_id=None, client_id=None, params=None))]
2028 #[expect(
2029 clippy::needless_pass_by_value,
2030 reason = "PyO3 owns extracted method arguments before Rust conversion"
2031 )]
2032 fn py_submit_order_list(
2033 &mut self,
2034 py: Python<'_>,
2035 order_list: Py<PyAny>,
2036 position_id: Option<PositionId>,
2037 client_id: Option<ClientId>,
2038 params: Option<Py<PyDict>>,
2039 ) -> PyResult<()> {
2040 let orders = py_order_list_to_orders(py, &order_list)?;
2041 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2042 match params {
2043 Some(dict) => from_pydict(py, &dict),
2044 None => Ok(None),
2045 }
2046 })?;
2047 let inner = self.inner_mut();
2048
2049 Strategy::submit_order_list(inner, orders, position_id, client_id, params_map)
2050 .map_err(to_pyruntime_err)
2051 }
2052
2053 #[pyo3(name = "modify_order")]
2054 #[pyo3(signature = (client_order_id, quantity=None, price=None, trigger_price=None, client_id=None, params=None))]
2055 fn py_modify_order(
2056 &mut self,
2057 client_order_id: ClientOrderId,
2058 quantity: Option<Quantity>,
2059 price: Option<Price>,
2060 trigger_price: Option<Price>,
2061 client_id: Option<ClientId>,
2062 params: Option<Py<PyDict>>,
2063 ) -> PyResult<()> {
2064 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2065 match params {
2066 Some(dict) => from_pydict(py, &dict),
2067 None => Ok(None),
2068 }
2069 })?;
2070 let inner = self.inner_mut();
2071
2072 Strategy::modify_order(
2073 inner,
2074 client_order_id,
2075 quantity,
2076 price,
2077 trigger_price,
2078 client_id,
2079 params_map,
2080 )
2081 .map_err(to_pyruntime_err)
2082 }
2083
2084 #[pyo3(name = "modify_orders")]
2085 #[pyo3(signature = (updates, client_id=None, params=None))]
2086 fn py_modify_orders(
2087 &mut self,
2088 updates: Vec<BatchModifyOrder>,
2089 client_id: Option<ClientId>,
2090 params: Option<Py<PyDict>>,
2091 ) -> PyResult<()> {
2092 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2093 match params {
2094 Some(dict) => from_pydict(py, &dict),
2095 None => Ok(None),
2096 }
2097 })?;
2098
2099 Strategy::modify_orders(self.inner_mut(), updates, client_id, params_map)
2100 .map_err(to_pyruntime_err)
2101 }
2102
2103 #[pyo3(name = "cancel_order")]
2104 #[pyo3(signature = (client_order_id, client_id=None, params=None))]
2105 fn py_cancel_order(
2106 &mut self,
2107 client_order_id: ClientOrderId,
2108 client_id: Option<ClientId>,
2109 params: Option<Py<PyDict>>,
2110 ) -> PyResult<()> {
2111 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2112 match params {
2113 Some(dict) => from_pydict(py, &dict),
2114 None => Ok(None),
2115 }
2116 })?;
2117 let inner = self.inner_mut();
2118
2119 Strategy::cancel_order(inner, client_order_id, client_id, params_map)
2120 .map_err(to_pyruntime_err)
2121 }
2122
2123 #[pyo3(name = "cancel_gtd_expiry")]
2125 #[pyo3(signature = (order))]
2126 fn py_cancel_gtd_expiry(&mut self, py: Python<'_>, order: Py<PyAny>) -> PyResult<()> {
2127 let order = pyobject_to_order_any(py, order)?;
2128
2129 Strategy::cancel_gtd_expiry(self.inner_mut(), &order.client_order_id());
2130 Ok(())
2131 }
2132
2133 #[pyo3(name = "cancel_orders")]
2134 #[pyo3(signature = (client_order_ids, client_id=None, params=None))]
2135 fn py_cancel_orders(
2136 &mut self,
2137 client_order_ids: Vec<ClientOrderId>,
2138 client_id: Option<ClientId>,
2139 params: Option<Py<PyDict>>,
2140 ) -> PyResult<()> {
2141 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2142 match params {
2143 Some(dict) => from_pydict(py, &dict),
2144 None => Ok(None),
2145 }
2146 })?;
2147
2148 Strategy::cancel_orders(self.inner_mut(), client_order_ids, client_id, params_map)
2149 .map_err(to_pyruntime_err)
2150 }
2151
2152 #[pyo3(name = "cancel_all_orders")]
2153 #[pyo3(signature = (instrument_id, order_side=None, client_id=None, strategy_only=true, params=None))]
2154 fn py_cancel_all_orders(
2155 &mut self,
2156 instrument_id: InstrumentId,
2157 order_side: Option<OrderSide>,
2158 client_id: Option<ClientId>,
2159 strategy_only: bool,
2160 params: Option<Py<PyDict>>,
2161 ) -> PyResult<()> {
2162 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2163 match params {
2164 Some(dict) => from_pydict(py, &dict),
2165 None => Ok(None),
2166 }
2167 })?;
2168 Strategy::cancel_all_orders(
2169 self.inner_mut(),
2170 instrument_id,
2171 order_side,
2172 client_id,
2173 strategy_only,
2174 params_map,
2175 )
2176 .map_err(to_pyruntime_err)
2177 }
2178
2179 #[pyo3(name = "close_position")]
2180 #[pyo3(signature = (position, client_id=None, tags=None, time_in_force=None, reduce_only=None, quote_quantity=None, params=None))]
2181 #[expect(clippy::too_many_arguments)]
2182 fn py_close_position(
2183 &mut self,
2184 position: &Position,
2185 client_id: Option<ClientId>,
2186 tags: Option<Vec<String>>,
2187 time_in_force: Option<TimeInForce>,
2188 reduce_only: Option<bool>,
2189 quote_quantity: Option<bool>,
2190 params: Option<Py<PyDict>>,
2191 ) -> PyResult<()> {
2192 let tags = tags.map(|t| t.into_iter().map(|s| Ustr::from(&s)).collect());
2193 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2194 match params {
2195 Some(dict) => from_pydict(py, &dict),
2196 None => Ok(None),
2197 }
2198 })?;
2199 Strategy::close_position(
2200 self.inner_mut(),
2201 position,
2202 client_id,
2203 tags,
2204 time_in_force,
2205 reduce_only,
2206 quote_quantity,
2207 params_map,
2208 )
2209 .map_err(to_pyruntime_err)
2210 }
2211
2212 #[pyo3(name = "close_all_positions")]
2213 #[pyo3(signature = (instrument_id, position_side=None, client_id=None, tags=None, time_in_force=None, reduce_only=None, quote_quantity=None, params=None))]
2214 #[expect(clippy::too_many_arguments)]
2215 fn py_close_all_positions(
2216 &mut self,
2217 instrument_id: InstrumentId,
2218 position_side: Option<PositionSide>,
2219 client_id: Option<ClientId>,
2220 tags: Option<Vec<String>>,
2221 time_in_force: Option<TimeInForce>,
2222 reduce_only: Option<bool>,
2223 quote_quantity: Option<bool>,
2224 params: Option<Py<PyDict>>,
2225 ) -> PyResult<()> {
2226 let tags = tags.map(|t| t.into_iter().map(|s| Ustr::from(&s)).collect());
2227 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2228 match params {
2229 Some(dict) => from_pydict(py, &dict),
2230 None => Ok(None),
2231 }
2232 })?;
2233 Strategy::close_all_positions(
2234 self.inner_mut(),
2235 instrument_id,
2236 position_side,
2237 client_id,
2238 tags,
2239 time_in_force,
2240 reduce_only,
2241 quote_quantity,
2242 params_map,
2243 )
2244 .map_err(to_pyruntime_err)
2245 }
2246
2247 #[pyo3(name = "query_account")]
2248 #[pyo3(signature = (account_id, client_id=None, params=None))]
2249 fn py_query_account(
2250 &mut self,
2251 py: Python<'_>,
2252 account_id: AccountId,
2253 client_id: Option<ClientId>,
2254 params: Option<Py<PyDict>>,
2255 ) -> PyResult<()> {
2256 let params_map = match params {
2257 Some(dict) => from_pydict(py, &dict)?,
2258 None => None,
2259 };
2260 Strategy::query_account(self.inner_mut(), account_id, client_id, params_map)
2261 .map_err(to_pyruntime_err)
2262 }
2263
2264 #[pyo3(name = "query_order")]
2265 #[pyo3(signature = (order, client_id=None, params=None))]
2266 fn py_query_order(
2267 &mut self,
2268 py: Python<'_>,
2269 order: Py<PyAny>,
2270 client_id: Option<ClientId>,
2271 params: Option<Py<PyDict>>,
2272 ) -> PyResult<()> {
2273 let order = pyobject_to_order_any(py, order)?;
2274 let params_map = match params {
2275 Some(dict) => from_pydict(py, &dict)?,
2276 None => None,
2277 };
2278 Strategy::query_order(self.inner_mut(), &order, client_id, params_map)
2279 .map_err(to_pyruntime_err)
2280 }
2281
2282 #[pyo3(name = "on_start")]
2283 fn py_on_start(&mut self) {}
2284
2285 #[pyo3(name = "on_stop")]
2286 fn py_on_stop(&mut self) {}
2287
2288 #[pyo3(name = "on_resume")]
2289 fn py_on_resume(&mut self) {}
2290
2291 #[pyo3(name = "on_reset")]
2292 fn py_on_reset(&mut self) {}
2293
2294 #[pyo3(name = "on_dispose")]
2295 fn py_on_dispose(&mut self) {}
2296
2297 #[pyo3(name = "on_degrade")]
2298 fn py_on_degrade(&mut self) {}
2299
2300 #[pyo3(name = "on_fault")]
2301 fn py_on_fault(&mut self) {}
2302
2303 #[pyo3(name = "on_save")]
2304 fn py_on_save(&self, py: Python<'_>) -> Py<PyDict> {
2305 PyDict::new(py).unbind()
2306 }
2307
2308 #[allow(unused_variables)]
2309 #[pyo3(name = "on_load")]
2310 fn py_on_load(&mut self, state: &Bound<'_, PyDict>) {}
2311
2312 #[allow(unused_variables, clippy::needless_pass_by_value)]
2313 #[pyo3(name = "on_time_event")]
2314 fn py_on_time_event(&mut self, event: TimeEvent) {}
2315
2316 #[allow(unused_variables, clippy::needless_pass_by_value)]
2317 #[pyo3(name = "on_data")]
2318 fn py_on_data(&mut self, data: Py<PyAny>) {}
2319
2320 #[allow(unused_variables)]
2321 #[pyo3(name = "on_signal")]
2322 fn py_on_signal(&mut self, signal: &Signal) {}
2323
2324 #[allow(unused_variables, clippy::needless_pass_by_value)]
2325 #[pyo3(name = "on_queue_state")]
2326 fn py_on_queue_state(&mut self, event: QueueStateChanged) {}
2327
2328 #[allow(unused_variables, clippy::needless_pass_by_value)]
2329 #[pyo3(name = "on_socket_state")]
2330 fn py_on_socket_state(&mut self, event: SocketStateChanged) {}
2331
2332 #[allow(unused_variables, clippy::needless_pass_by_value)]
2333 #[pyo3(name = "on_instrument")]
2334 fn py_on_instrument(&mut self, instrument: Py<PyAny>) {}
2335
2336 #[allow(unused_variables)]
2337 #[pyo3(name = "on_quote")]
2338 fn py_on_quote(&mut self, quote: QuoteTick) {}
2339
2340 #[allow(unused_variables)]
2341 #[pyo3(name = "on_trade")]
2342 fn py_on_trade(&mut self, trade: TradeTick) {}
2343
2344 #[allow(unused_variables)]
2345 #[pyo3(name = "on_bar")]
2346 fn py_on_bar(&mut self, bar: Bar) {}
2347
2348 #[allow(unused_variables, clippy::needless_pass_by_value)]
2349 #[pyo3(name = "on_book_deltas")]
2350 fn py_on_book_deltas(&mut self, deltas: OrderBookDeltas) {}
2351
2352 #[allow(unused_variables)]
2353 #[pyo3(name = "on_book_depth")]
2354 fn py_on_book_depth(&mut self, depth: &OrderBookDepth) {}
2355
2356 #[allow(unused_variables)]
2357 #[pyo3(name = "on_book")]
2358 fn py_on_book(&mut self, book: &OrderBook) {}
2359
2360 #[allow(unused_variables)]
2361 #[pyo3(name = "on_mark_price")]
2362 fn py_on_mark_price(&mut self, mark_price: MarkPriceUpdate) {}
2363
2364 #[allow(unused_variables)]
2365 #[pyo3(name = "on_index_price")]
2366 fn py_on_index_price(&mut self, index_price: IndexPriceUpdate) {}
2367
2368 #[allow(unused_variables)]
2369 #[pyo3(name = "on_funding_rate")]
2370 fn py_on_funding_rate(&mut self, funding_rate: FundingRateUpdate) {}
2371
2372 #[allow(unused_variables)]
2373 #[pyo3(name = "on_instrument_status")]
2374 fn py_on_instrument_status(&mut self, status: InstrumentStatus) {}
2375
2376 #[allow(unused_variables)]
2377 #[pyo3(name = "on_instrument_close")]
2378 fn py_on_instrument_close(&mut self, close: InstrumentClose) {}
2379
2380 #[allow(unused_variables)]
2381 #[pyo3(name = "on_option_greeks")]
2382 fn py_on_option_greeks(&mut self, greeks: OptionGreeks) {}
2383
2384 #[allow(unused_variables, clippy::needless_pass_by_value)]
2385 #[pyo3(name = "on_option_chain")]
2386 fn py_on_option_chain(&mut self, slice: OptionChainSlice) {}
2387
2388 #[pyo3(name = "on_market_exit")]
2389 fn py_on_market_exit(&mut self) {}
2390
2391 #[pyo3(name = "post_market_exit")]
2392 fn py_post_market_exit(&mut self) {}
2393
2394 #[allow(unused_variables, clippy::needless_pass_by_value)]
2395 #[pyo3(name = "on_order_initialized")]
2396 fn py_on_order_initialized(&mut self, event: OrderInitialized) {}
2397
2398 #[allow(unused_variables, clippy::needless_pass_by_value)]
2399 #[pyo3(name = "on_order_event")]
2400 fn py_on_order_event(&mut self, event: Py<PyAny>) {}
2401
2402 #[allow(unused_variables)]
2403 #[pyo3(name = "on_order_denied")]
2404 fn py_on_order_denied(&mut self, event: OrderDenied) {}
2405
2406 #[allow(unused_variables)]
2407 #[pyo3(name = "on_order_emulated")]
2408 fn py_on_order_emulated(&mut self, event: OrderEmulated) {}
2409
2410 #[allow(unused_variables)]
2411 #[pyo3(name = "on_order_released")]
2412 fn py_on_order_released(&mut self, event: OrderReleased) {}
2413
2414 #[allow(unused_variables)]
2415 #[pyo3(name = "on_order_submitted")]
2416 fn py_on_order_submitted(&mut self, event: OrderSubmitted) {}
2417
2418 #[allow(unused_variables)]
2419 #[pyo3(name = "on_order_rejected")]
2420 fn py_on_order_rejected(&mut self, event: OrderRejected) {}
2421
2422 #[allow(unused_variables)]
2423 #[pyo3(name = "on_order_accepted")]
2424 fn py_on_order_accepted(&mut self, event: OrderAccepted) {}
2425
2426 #[allow(unused_variables)]
2427 #[pyo3(name = "on_order_expired")]
2428 fn py_on_order_expired(&mut self, event: OrderExpired) {}
2429
2430 #[allow(unused_variables)]
2431 #[pyo3(name = "on_order_triggered")]
2432 fn py_on_order_triggered(&mut self, event: OrderTriggered) {}
2433
2434 #[allow(unused_variables)]
2435 #[pyo3(name = "on_order_pending_update")]
2436 fn py_on_order_pending_update(&mut self, event: OrderPendingUpdate) {}
2437
2438 #[allow(unused_variables)]
2439 #[pyo3(name = "on_order_pending_cancel")]
2440 fn py_on_order_pending_cancel(&mut self, event: OrderPendingCancel) {}
2441
2442 #[allow(unused_variables)]
2443 #[pyo3(name = "on_order_modify_rejected")]
2444 fn py_on_order_modify_rejected(&mut self, event: OrderModifyRejected) {}
2445
2446 #[allow(unused_variables)]
2447 #[pyo3(name = "on_order_cancel_rejected")]
2448 fn py_on_order_cancel_rejected(&mut self, event: OrderCancelRejected) {}
2449
2450 #[allow(unused_variables)]
2451 #[pyo3(name = "on_order_updated")]
2452 fn py_on_order_updated(&mut self, event: OrderUpdated) {}
2453
2454 #[allow(unused_variables)]
2455 #[pyo3(name = "on_order_canceled")]
2456 fn py_on_order_canceled(&mut self, event: OrderCanceled) {}
2457
2458 #[allow(unused_variables, clippy::needless_pass_by_value)]
2459 #[pyo3(name = "on_order_filled")]
2460 fn py_on_order_filled(&mut self, event: OrderFilled) {}
2461
2462 #[allow(unused_variables, clippy::needless_pass_by_value)]
2463 #[pyo3(name = "on_order_fill_voided")]
2464 fn py_on_order_fill_voided(&mut self, event: OrderFillVoided) {}
2465
2466 #[allow(unused_variables, clippy::needless_pass_by_value)]
2467 #[pyo3(name = "on_position_opened")]
2468 fn py_on_position_opened(&mut self, event: PositionOpened) {}
2469
2470 #[allow(unused_variables, clippy::needless_pass_by_value)]
2471 #[pyo3(name = "on_position_event")]
2472 fn py_on_position_event(&mut self, event: Py<PyAny>) {}
2473
2474 #[allow(unused_variables, clippy::needless_pass_by_value)]
2475 #[pyo3(name = "on_position_changed")]
2476 fn py_on_position_changed(&mut self, event: PositionChanged) {}
2477
2478 #[allow(unused_variables, clippy::needless_pass_by_value)]
2479 #[pyo3(name = "on_position_closed")]
2480 fn py_on_position_closed(&mut self, event: PositionClosed) {}
2481
2482 #[allow(unused_variables, clippy::needless_pass_by_value)]
2483 #[pyo3(name = "on_historical_data")]
2484 fn py_on_historical_data(&mut self, data: Py<PyAny>) {
2485 }
2487
2488 #[allow(unused_variables, clippy::needless_pass_by_value)]
2489 #[pyo3(name = "on_historical_book_deltas")]
2490 fn py_on_historical_book_deltas(&mut self, deltas: Vec<OrderBookDelta>) {}
2491
2492 #[allow(unused_variables, clippy::needless_pass_by_value)]
2493 #[pyo3(name = "on_historical_book_depth")]
2494 fn py_on_historical_book_depth(&mut self, depths: Vec<OrderBookDepth>) {}
2495
2496 #[allow(unused_variables, clippy::needless_pass_by_value)]
2497 #[pyo3(name = "on_historical_quotes")]
2498 fn py_on_historical_quotes(&mut self, quotes: Vec<QuoteTick>) {
2499 }
2501
2502 #[allow(unused_variables, clippy::needless_pass_by_value)]
2503 #[pyo3(name = "on_historical_trades")]
2504 fn py_on_historical_trades(&mut self, trades: Vec<TradeTick>) {
2505 }
2507
2508 #[allow(unused_variables, clippy::needless_pass_by_value)]
2509 #[pyo3(name = "on_historical_funding_rates")]
2510 fn py_on_historical_funding_rates(&mut self, funding_rates: Vec<FundingRateUpdate>) {
2511 }
2513
2514 #[allow(unused_variables, clippy::needless_pass_by_value)]
2515 #[pyo3(name = "on_historical_bars")]
2516 fn py_on_historical_bars(&mut self, bars: Vec<Bar>) {
2517 }
2519
2520 #[allow(unused_variables, clippy::needless_pass_by_value)]
2521 #[pyo3(name = "on_historical_mark_prices")]
2522 fn py_on_historical_mark_prices(&mut self, mark_prices: Vec<MarkPriceUpdate>) {
2523 }
2525
2526 #[allow(unused_variables, clippy::needless_pass_by_value)]
2527 #[pyo3(name = "on_historical_index_prices")]
2528 fn py_on_historical_index_prices(&mut self, index_prices: Vec<IndexPriceUpdate>) {
2529 }
2531
2532 #[pyo3(name = "subscribe_data")]
2533 #[pyo3(signature = (data_type, client_id=None, params=None))]
2534 fn py_subscribe_data(
2535 &mut self,
2536 data_type: DataType,
2537 client_id: Option<ClientId>,
2538 params: Option<Py<PyDict>>,
2539 ) -> PyResult<()> {
2540 self.ensure_registered()?;
2541 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2542 match params {
2543 Some(dict) => from_pydict(py, &dict),
2544 None => Ok(None),
2545 }
2546 })?;
2547 DataActor::subscribe_data(self.inner_mut(), data_type, client_id, params_map);
2548 Ok(())
2549 }
2550
2551 #[pyo3(name = "subscribe_signal")]
2552 #[pyo3(signature = (name="", priority=None))]
2553 fn py_subscribe_signal(&mut self, name: &str, priority: Option<u32>) -> PyResult<()> {
2554 self.ensure_registered()?;
2555 DataActor::subscribe_signal(self.inner_mut(), name, priority);
2556 Ok(())
2557 }
2558
2559 #[pyo3(name = "subscribe_queue_state")]
2560 #[pyo3(signature = (channel=None, priority=None))]
2561 fn py_subscribe_queue_state(
2562 &mut self,
2563 channel: Option<SystemChannel>,
2564 priority: Option<u32>,
2565 ) -> PyResult<()> {
2566 self.ensure_registered()?;
2567 DataActor::subscribe_queue_state(self.inner_mut(), channel, priority);
2568 Ok(())
2569 }
2570
2571 #[pyo3(name = "subscribe_socket_state")]
2572 #[pyo3(signature = (client_id=None, endpoint=None, priority=None))]
2573 fn py_subscribe_socket_state(
2574 &mut self,
2575 client_id: Option<ClientId>,
2576 endpoint: Option<&str>,
2577 priority: Option<u32>,
2578 ) -> PyResult<()> {
2579 self.ensure_registered()?;
2580 DataActor::subscribe_socket_state(self.inner_mut(), client_id, endpoint, priority);
2581 Ok(())
2582 }
2583
2584 #[pyo3(name = "subscribe_instruments")]
2585 #[pyo3(signature = (venue, client_id=None, params=None))]
2586 fn py_subscribe_instruments(
2587 &mut self,
2588 venue: Venue,
2589 client_id: Option<ClientId>,
2590 params: Option<Py<PyDict>>,
2591 ) -> PyResult<()> {
2592 self.ensure_registered()?;
2593 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2594 match params {
2595 Some(dict) => from_pydict(py, &dict),
2596 None => Ok(None),
2597 }
2598 })?;
2599 DataActor::subscribe_instruments(self.inner_mut(), venue, client_id, params_map);
2600 Ok(())
2601 }
2602
2603 #[pyo3(name = "subscribe_instrument")]
2604 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2605 fn py_subscribe_instrument(
2606 &mut self,
2607 instrument_id: InstrumentId,
2608 client_id: Option<ClientId>,
2609 params: Option<Py<PyDict>>,
2610 ) -> PyResult<()> {
2611 self.ensure_registered()?;
2612 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2613 match params {
2614 Some(dict) => from_pydict(py, &dict),
2615 None => Ok(None),
2616 }
2617 })?;
2618 DataActor::subscribe_instrument(self.inner_mut(), instrument_id, client_id, params_map);
2619 Ok(())
2620 }
2621
2622 #[pyo3(name = "subscribe_book_deltas")]
2623 #[pyo3(signature = (instrument_id, book_type, depth=None, client_id=None, managed=false, params=None))]
2624 fn py_subscribe_book_deltas(
2625 &mut self,
2626 instrument_id: InstrumentId,
2627 book_type: BookType,
2628 depth: Option<usize>,
2629 client_id: Option<ClientId>,
2630 managed: bool,
2631 params: Option<Py<PyDict>>,
2632 ) -> PyResult<()> {
2633 self.ensure_registered()?;
2634 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2635 match params {
2636 Some(dict) => from_pydict(py, &dict),
2637 None => Ok(None),
2638 }
2639 })?;
2640 let depth = depth.and_then(NonZeroUsize::new);
2641 DataActor::subscribe_book_deltas(
2642 self.inner_mut(),
2643 instrument_id,
2644 book_type,
2645 depth,
2646 client_id,
2647 managed,
2648 params_map,
2649 );
2650 Ok(())
2651 }
2652
2653 #[pyo3(name = "subscribe_book_depth")]
2654 #[pyo3(signature = (instrument_id, book_type, depth=None, client_id=None, managed=false, params=None))]
2655 fn py_subscribe_book_depth(
2656 &mut self,
2657 instrument_id: InstrumentId,
2658 book_type: BookType,
2659 depth: Option<usize>,
2660 client_id: Option<ClientId>,
2661 managed: bool,
2662 params: Option<Py<PyDict>>,
2663 ) -> PyResult<()> {
2664 self.ensure_registered()?;
2665
2666 let depth = depth
2667 .map(|value| {
2668 NonZeroUsize::new(value).ok_or_else(|| to_pyvalue_err("depth must be positive"))
2669 })
2670 .transpose()?;
2671
2672 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2673 match params {
2674 Some(dict) => from_pydict(py, &dict),
2675 None => Ok(None),
2676 }
2677 })?;
2678 DataActor::subscribe_book_depth(
2679 self.inner_mut(),
2680 instrument_id,
2681 book_type,
2682 depth,
2683 client_id,
2684 managed,
2685 params_map,
2686 );
2687 Ok(())
2688 }
2689
2690 #[pyo3(name = "subscribe_book_at_interval")]
2691 #[pyo3(signature = (instrument_id, book_type, interval_ms, depth=None, client_id=None, params=None))]
2692 fn py_subscribe_book_at_interval(
2693 &mut self,
2694 instrument_id: InstrumentId,
2695 book_type: BookType,
2696 interval_ms: usize,
2697 depth: Option<usize>,
2698 client_id: Option<ClientId>,
2699 params: Option<Py<PyDict>>,
2700 ) -> PyResult<()> {
2701 let interval_ms = NonZeroUsize::new(interval_ms)
2702 .ok_or_else(|| to_pyvalue_err("interval_ms must be > 0"))?;
2703
2704 self.ensure_registered()?;
2705 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2706 match params {
2707 Some(dict) => from_pydict(py, &dict),
2708 None => Ok(None),
2709 }
2710 })?;
2711 let depth = depth.and_then(NonZeroUsize::new);
2712 DataActor::subscribe_book_at_interval(
2713 self.inner_mut(),
2714 instrument_id,
2715 book_type,
2716 depth,
2717 interval_ms,
2718 client_id,
2719 params_map,
2720 );
2721 Ok(())
2722 }
2723
2724 #[pyo3(name = "subscribe_quotes")]
2725 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2726 fn py_subscribe_quotes(
2727 &mut self,
2728 instrument_id: InstrumentId,
2729 client_id: Option<ClientId>,
2730 params: Option<Py<PyDict>>,
2731 ) -> PyResult<()> {
2732 self.ensure_registered()?;
2733 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2734 match params {
2735 Some(dict) => from_pydict(py, &dict),
2736 None => Ok(None),
2737 }
2738 })?;
2739 DataActor::subscribe_quotes(self.inner_mut(), instrument_id, client_id, params_map);
2740 Ok(())
2741 }
2742
2743 #[pyo3(name = "subscribe_trades")]
2744 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2745 fn py_subscribe_trades(
2746 &mut self,
2747 instrument_id: InstrumentId,
2748 client_id: Option<ClientId>,
2749 params: Option<Py<PyDict>>,
2750 ) -> PyResult<()> {
2751 self.ensure_registered()?;
2752 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2753 match params {
2754 Some(dict) => from_pydict(py, &dict),
2755 None => Ok(None),
2756 }
2757 })?;
2758 DataActor::subscribe_trades(self.inner_mut(), instrument_id, client_id, params_map);
2759 Ok(())
2760 }
2761
2762 #[pyo3(name = "subscribe_bars")]
2763 #[pyo3(signature = (bar_type, client_id=None, params=None))]
2764 fn py_subscribe_bars(
2765 &mut self,
2766 bar_type: BarType,
2767 client_id: Option<ClientId>,
2768 params: Option<Py<PyDict>>,
2769 ) -> PyResult<()> {
2770 self.ensure_registered()?;
2771 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2772 match params {
2773 Some(dict) => from_pydict(py, &dict),
2774 None => Ok(None),
2775 }
2776 })?;
2777 DataActor::subscribe_bars(self.inner_mut(), bar_type, client_id, params_map);
2778 Ok(())
2779 }
2780
2781 #[pyo3(name = "subscribe_mark_prices")]
2782 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2783 fn py_subscribe_mark_prices(
2784 &mut self,
2785 instrument_id: InstrumentId,
2786 client_id: Option<ClientId>,
2787 params: Option<Py<PyDict>>,
2788 ) -> PyResult<()> {
2789 self.ensure_registered()?;
2790 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2791 match params {
2792 Some(dict) => from_pydict(py, &dict),
2793 None => Ok(None),
2794 }
2795 })?;
2796 DataActor::subscribe_mark_prices(self.inner_mut(), instrument_id, client_id, params_map);
2797 Ok(())
2798 }
2799
2800 #[pyo3(name = "subscribe_index_prices")]
2801 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2802 fn py_subscribe_index_prices(
2803 &mut self,
2804 instrument_id: InstrumentId,
2805 client_id: Option<ClientId>,
2806 params: Option<Py<PyDict>>,
2807 ) -> PyResult<()> {
2808 self.ensure_registered()?;
2809 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2810 match params {
2811 Some(dict) => from_pydict(py, &dict),
2812 None => Ok(None),
2813 }
2814 })?;
2815 DataActor::subscribe_index_prices(self.inner_mut(), instrument_id, client_id, params_map);
2816 Ok(())
2817 }
2818
2819 #[pyo3(name = "subscribe_funding_rates")]
2820 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2821 fn py_subscribe_funding_rates(
2822 &mut self,
2823 instrument_id: InstrumentId,
2824 client_id: Option<ClientId>,
2825 params: Option<Py<PyDict>>,
2826 ) -> PyResult<()> {
2827 self.ensure_registered()?;
2828 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2829 match params {
2830 Some(dict) => from_pydict(py, &dict),
2831 None => Ok(None),
2832 }
2833 })?;
2834 DataActor::subscribe_funding_rates(self.inner_mut(), instrument_id, client_id, params_map);
2835 Ok(())
2836 }
2837
2838 #[pyo3(name = "subscribe_option_greeks")]
2839 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2840 fn py_subscribe_option_greeks(
2841 &mut self,
2842 instrument_id: InstrumentId,
2843 client_id: Option<ClientId>,
2844 params: Option<Py<PyDict>>,
2845 ) -> PyResult<()> {
2846 self.ensure_registered()?;
2847 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2848 match params {
2849 Some(dict) => from_pydict(py, &dict),
2850 None => Ok(None),
2851 }
2852 })?;
2853 DataActor::subscribe_option_greeks(self.inner_mut(), instrument_id, client_id, params_map);
2854 Ok(())
2855 }
2856
2857 #[pyo3(name = "subscribe_instrument_status")]
2858 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2859 fn py_subscribe_instrument_status(
2860 &mut self,
2861 instrument_id: InstrumentId,
2862 client_id: Option<ClientId>,
2863 params: Option<Py<PyDict>>,
2864 ) -> PyResult<()> {
2865 self.ensure_registered()?;
2866 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2867 match params {
2868 Some(dict) => from_pydict(py, &dict),
2869 None => Ok(None),
2870 }
2871 })?;
2872 DataActor::subscribe_instrument_status(
2873 self.inner_mut(),
2874 instrument_id,
2875 client_id,
2876 params_map,
2877 );
2878 Ok(())
2879 }
2880
2881 #[pyo3(name = "subscribe_instrument_close")]
2882 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2883 fn py_subscribe_instrument_close(
2884 &mut self,
2885 instrument_id: InstrumentId,
2886 client_id: Option<ClientId>,
2887 params: Option<Py<PyDict>>,
2888 ) -> PyResult<()> {
2889 self.ensure_registered()?;
2890 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2891 match params {
2892 Some(dict) => from_pydict(py, &dict),
2893 None => Ok(None),
2894 }
2895 })?;
2896 DataActor::subscribe_instrument_close(
2897 self.inner_mut(),
2898 instrument_id,
2899 client_id,
2900 params_map,
2901 );
2902 Ok(())
2903 }
2904
2905 #[pyo3(name = "subscribe_option_chain")]
2906 #[pyo3(signature = (series_id, strike_range, snapshot_interval_ms=None, client_id=None, params=None))]
2907 fn py_subscribe_option_chain(
2908 &mut self,
2909 py: Python<'_>,
2910 series_id: OptionSeriesId,
2911 strike_range: PyStrikeRange,
2912 snapshot_interval_ms: Option<u64>,
2913 client_id: Option<ClientId>,
2914 params: Option<Py<PyDict>>,
2915 ) -> PyResult<()> {
2916 self.ensure_registered()?;
2917 let params_map = match params {
2918 Some(dict) => from_pydict(py, &dict)?,
2919 None => None,
2920 };
2921 DataActor::subscribe_option_chain(
2922 self.inner_mut(),
2923 series_id,
2924 strike_range.inner,
2925 snapshot_interval_ms,
2926 client_id,
2927 params_map,
2928 );
2929 Ok(())
2930 }
2931
2932 #[pyo3(name = "unsubscribe_data")]
2933 #[pyo3(signature = (data_type, client_id=None, params=None))]
2934 fn py_unsubscribe_data(
2935 &mut self,
2936 data_type: DataType,
2937 client_id: Option<ClientId>,
2938 params: Option<Py<PyDict>>,
2939 ) -> PyResult<()> {
2940 self.ensure_registered()?;
2941 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2942 match params {
2943 Some(dict) => from_pydict(py, &dict),
2944 None => Ok(None),
2945 }
2946 })?;
2947 DataActor::unsubscribe_data(self.inner_mut(), data_type, client_id, params_map);
2948 Ok(())
2949 }
2950
2951 #[pyo3(name = "unsubscribe_signal")]
2952 fn py_unsubscribe_signal(&mut self, name: &str) -> PyResult<()> {
2953 self.ensure_registered()?;
2954 DataActor::unsubscribe_signal(self.inner_mut(), name);
2955 Ok(())
2956 }
2957
2958 #[pyo3(name = "unsubscribe_queue_state")]
2959 #[pyo3(signature = (channel=None))]
2960 fn py_unsubscribe_queue_state(&mut self, channel: Option<SystemChannel>) -> PyResult<()> {
2961 self.ensure_registered()?;
2962 DataActor::unsubscribe_queue_state(self.inner_mut(), channel);
2963 Ok(())
2964 }
2965
2966 #[pyo3(name = "unsubscribe_socket_state")]
2967 #[pyo3(signature = (client_id=None, endpoint=None))]
2968 fn py_unsubscribe_socket_state(
2969 &mut self,
2970 client_id: Option<ClientId>,
2971 endpoint: Option<&str>,
2972 ) -> PyResult<()> {
2973 self.ensure_registered()?;
2974 DataActor::unsubscribe_socket_state(self.inner_mut(), client_id, endpoint);
2975 Ok(())
2976 }
2977
2978 #[pyo3(name = "unsubscribe_instruments")]
2979 #[pyo3(signature = (venue, client_id=None, params=None))]
2980 fn py_unsubscribe_instruments(
2981 &mut self,
2982 venue: Venue,
2983 client_id: Option<ClientId>,
2984 params: Option<Py<PyDict>>,
2985 ) -> PyResult<()> {
2986 self.ensure_registered()?;
2987 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
2988 match params {
2989 Some(dict) => from_pydict(py, &dict),
2990 None => Ok(None),
2991 }
2992 })?;
2993 DataActor::unsubscribe_instruments(self.inner_mut(), venue, client_id, params_map);
2994 Ok(())
2995 }
2996
2997 #[pyo3(name = "unsubscribe_instrument")]
2998 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
2999 fn py_unsubscribe_instrument(
3000 &mut self,
3001 instrument_id: InstrumentId,
3002 client_id: Option<ClientId>,
3003 params: Option<Py<PyDict>>,
3004 ) -> PyResult<()> {
3005 self.ensure_registered()?;
3006 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3007 match params {
3008 Some(dict) => from_pydict(py, &dict),
3009 None => Ok(None),
3010 }
3011 })?;
3012 DataActor::unsubscribe_instrument(self.inner_mut(), instrument_id, client_id, params_map);
3013 Ok(())
3014 }
3015
3016 #[pyo3(name = "unsubscribe_book_deltas")]
3017 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3018 fn py_unsubscribe_book_deltas(
3019 &mut self,
3020 instrument_id: InstrumentId,
3021 client_id: Option<ClientId>,
3022 params: Option<Py<PyDict>>,
3023 ) -> PyResult<()> {
3024 self.ensure_registered()?;
3025 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3026 match params {
3027 Some(dict) => from_pydict(py, &dict),
3028 None => Ok(None),
3029 }
3030 })?;
3031 DataActor::unsubscribe_book_deltas(self.inner_mut(), instrument_id, client_id, params_map);
3032 Ok(())
3033 }
3034
3035 #[pyo3(name = "unsubscribe_book_depth")]
3036 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3037 fn py_unsubscribe_book_depth(
3038 &mut self,
3039 instrument_id: InstrumentId,
3040 client_id: Option<ClientId>,
3041 params: Option<Py<PyDict>>,
3042 ) -> PyResult<()> {
3043 self.ensure_registered()?;
3044 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3045 match params {
3046 Some(dict) => from_pydict(py, &dict),
3047 None => Ok(None),
3048 }
3049 })?;
3050 DataActor::unsubscribe_book_depth(self.inner_mut(), instrument_id, client_id, params_map);
3051 Ok(())
3052 }
3053
3054 #[pyo3(name = "unsubscribe_book_at_interval")]
3055 #[pyo3(signature = (instrument_id, interval_ms, client_id=None, params=None))]
3056 fn py_unsubscribe_book_at_interval(
3057 &mut self,
3058 instrument_id: InstrumentId,
3059 interval_ms: usize,
3060 client_id: Option<ClientId>,
3061 params: Option<Py<PyDict>>,
3062 ) -> PyResult<()> {
3063 let interval_ms = NonZeroUsize::new(interval_ms)
3064 .ok_or_else(|| to_pyvalue_err("interval_ms must be > 0"))?;
3065
3066 self.ensure_registered()?;
3067 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3068 match params {
3069 Some(dict) => from_pydict(py, &dict),
3070 None => Ok(None),
3071 }
3072 })?;
3073 DataActor::unsubscribe_book_at_interval(
3074 self.inner_mut(),
3075 instrument_id,
3076 interval_ms,
3077 client_id,
3078 params_map,
3079 );
3080 Ok(())
3081 }
3082
3083 #[pyo3(name = "unsubscribe_quotes")]
3084 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3085 fn py_unsubscribe_quotes(
3086 &mut self,
3087 instrument_id: InstrumentId,
3088 client_id: Option<ClientId>,
3089 params: Option<Py<PyDict>>,
3090 ) -> PyResult<()> {
3091 self.ensure_registered()?;
3092 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3093 match params {
3094 Some(dict) => from_pydict(py, &dict),
3095 None => Ok(None),
3096 }
3097 })?;
3098 DataActor::unsubscribe_quotes(self.inner_mut(), instrument_id, client_id, params_map);
3099 Ok(())
3100 }
3101
3102 #[pyo3(name = "unsubscribe_trades")]
3103 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3104 fn py_unsubscribe_trades(
3105 &mut self,
3106 instrument_id: InstrumentId,
3107 client_id: Option<ClientId>,
3108 params: Option<Py<PyDict>>,
3109 ) -> PyResult<()> {
3110 self.ensure_registered()?;
3111 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3112 match params {
3113 Some(dict) => from_pydict(py, &dict),
3114 None => Ok(None),
3115 }
3116 })?;
3117 DataActor::unsubscribe_trades(self.inner_mut(), instrument_id, client_id, params_map);
3118 Ok(())
3119 }
3120
3121 #[pyo3(name = "unsubscribe_bars")]
3122 #[pyo3(signature = (bar_type, client_id=None, params=None))]
3123 fn py_unsubscribe_bars(
3124 &mut self,
3125 bar_type: BarType,
3126 client_id: Option<ClientId>,
3127 params: Option<Py<PyDict>>,
3128 ) -> PyResult<()> {
3129 self.ensure_registered()?;
3130 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3131 match params {
3132 Some(dict) => from_pydict(py, &dict),
3133 None => Ok(None),
3134 }
3135 })?;
3136 DataActor::unsubscribe_bars(self.inner_mut(), bar_type, client_id, params_map);
3137 Ok(())
3138 }
3139
3140 #[pyo3(name = "unsubscribe_mark_prices")]
3141 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3142 fn py_unsubscribe_mark_prices(
3143 &mut self,
3144 instrument_id: InstrumentId,
3145 client_id: Option<ClientId>,
3146 params: Option<Py<PyDict>>,
3147 ) -> PyResult<()> {
3148 self.ensure_registered()?;
3149 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3150 match params {
3151 Some(dict) => from_pydict(py, &dict),
3152 None => Ok(None),
3153 }
3154 })?;
3155 DataActor::unsubscribe_mark_prices(self.inner_mut(), instrument_id, client_id, params_map);
3156 Ok(())
3157 }
3158
3159 #[pyo3(name = "unsubscribe_index_prices")]
3160 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3161 fn py_unsubscribe_index_prices(
3162 &mut self,
3163 instrument_id: InstrumentId,
3164 client_id: Option<ClientId>,
3165 params: Option<Py<PyDict>>,
3166 ) -> PyResult<()> {
3167 self.ensure_registered()?;
3168 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3169 match params {
3170 Some(dict) => from_pydict(py, &dict),
3171 None => Ok(None),
3172 }
3173 })?;
3174 DataActor::unsubscribe_index_prices(self.inner_mut(), instrument_id, client_id, params_map);
3175 Ok(())
3176 }
3177
3178 #[pyo3(name = "unsubscribe_funding_rates")]
3179 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3180 fn py_unsubscribe_funding_rates(
3181 &mut self,
3182 instrument_id: InstrumentId,
3183 client_id: Option<ClientId>,
3184 params: Option<Py<PyDict>>,
3185 ) -> PyResult<()> {
3186 self.ensure_registered()?;
3187 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3188 match params {
3189 Some(dict) => from_pydict(py, &dict),
3190 None => Ok(None),
3191 }
3192 })?;
3193 DataActor::unsubscribe_funding_rates(
3194 self.inner_mut(),
3195 instrument_id,
3196 client_id,
3197 params_map,
3198 );
3199 Ok(())
3200 }
3201
3202 #[pyo3(name = "unsubscribe_option_greeks")]
3203 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3204 fn py_unsubscribe_option_greeks(
3205 &mut self,
3206 instrument_id: InstrumentId,
3207 client_id: Option<ClientId>,
3208 params: Option<Py<PyDict>>,
3209 ) -> PyResult<()> {
3210 self.ensure_registered()?;
3211 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3212 match params {
3213 Some(dict) => from_pydict(py, &dict),
3214 None => Ok(None),
3215 }
3216 })?;
3217 DataActor::unsubscribe_option_greeks(
3218 self.inner_mut(),
3219 instrument_id,
3220 client_id,
3221 params_map,
3222 );
3223 Ok(())
3224 }
3225
3226 #[pyo3(name = "unsubscribe_instrument_status")]
3227 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3228 fn py_unsubscribe_instrument_status(
3229 &mut self,
3230 instrument_id: InstrumentId,
3231 client_id: Option<ClientId>,
3232 params: Option<Py<PyDict>>,
3233 ) -> PyResult<()> {
3234 self.ensure_registered()?;
3235 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3236 match params {
3237 Some(dict) => from_pydict(py, &dict),
3238 None => Ok(None),
3239 }
3240 })?;
3241 DataActor::unsubscribe_instrument_status(
3242 self.inner_mut(),
3243 instrument_id,
3244 client_id,
3245 params_map,
3246 );
3247 Ok(())
3248 }
3249
3250 #[pyo3(name = "unsubscribe_instrument_close")]
3251 #[pyo3(signature = (instrument_id, client_id=None, params=None))]
3252 fn py_unsubscribe_instrument_close(
3253 &mut self,
3254 instrument_id: InstrumentId,
3255 client_id: Option<ClientId>,
3256 params: Option<Py<PyDict>>,
3257 ) -> PyResult<()> {
3258 self.ensure_registered()?;
3259 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3260 match params {
3261 Some(dict) => from_pydict(py, &dict),
3262 None => Ok(None),
3263 }
3264 })?;
3265 DataActor::unsubscribe_instrument_close(
3266 self.inner_mut(),
3267 instrument_id,
3268 client_id,
3269 params_map,
3270 );
3271 Ok(())
3272 }
3273
3274 #[pyo3(name = "unsubscribe_option_chain")]
3275 #[pyo3(signature = (series_id, client_id=None))]
3276 fn py_unsubscribe_option_chain(
3277 &mut self,
3278 series_id: OptionSeriesId,
3279 client_id: Option<ClientId>,
3280 ) -> PyResult<()> {
3281 self.ensure_registered()?;
3282 DataActor::unsubscribe_option_chain(self.inner_mut(), series_id, client_id);
3283 Ok(())
3284 }
3285
3286 #[pyo3(name = "request_data")]
3287 #[pyo3(signature = (data_type, client_id, start=None, end=None, limit=None, params=None))]
3288 fn py_request_data(
3289 &mut self,
3290 data_type: DataType,
3291 client_id: ClientId,
3292 start: Option<Timestamp>,
3293 end: Option<Timestamp>,
3294 limit: Option<usize>,
3295 params: Option<Py<PyDict>>,
3296 ) -> PyResult<String> {
3297 self.ensure_registered_for_data()?;
3298 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3299 match params {
3300 Some(dict) => from_pydict(py, &dict),
3301 None => Ok(None),
3302 }
3303 })?;
3304 let limit = limit.and_then(NonZeroUsize::new);
3305 let request_id = DataActor::request_data(
3306 self.inner_mut(),
3307 data_type,
3308 client_id,
3309 start,
3310 end,
3311 limit,
3312 params_map,
3313 )
3314 .map_err(to_pyvalue_err)?;
3315 Ok(request_id.to_string())
3316 }
3317
3318 #[pyo3(name = "request_instrument")]
3319 #[pyo3(signature = (instrument_id, start=None, end=None, client_id=None, params=None))]
3320 fn py_request_instrument(
3321 &mut self,
3322 instrument_id: InstrumentId,
3323 start: Option<Timestamp>,
3324 end: Option<Timestamp>,
3325 client_id: Option<ClientId>,
3326 params: Option<Py<PyDict>>,
3327 ) -> PyResult<String> {
3328 self.ensure_registered_for_data()?;
3329 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3330 match params {
3331 Some(dict) => from_pydict(py, &dict),
3332 None => Ok(None),
3333 }
3334 })?;
3335 let request_id = DataActor::request_instrument(
3336 self.inner_mut(),
3337 instrument_id,
3338 start,
3339 end,
3340 client_id,
3341 params_map,
3342 )
3343 .map_err(to_pyvalue_err)?;
3344 Ok(request_id.to_string())
3345 }
3346
3347 #[pyo3(name = "request_instruments")]
3348 #[pyo3(signature = (venue=None, start=None, end=None, client_id=None, params=None))]
3349 fn py_request_instruments(
3350 &mut self,
3351 venue: Option<Venue>,
3352 start: Option<Timestamp>,
3353 end: Option<Timestamp>,
3354 client_id: Option<ClientId>,
3355 params: Option<Py<PyDict>>,
3356 ) -> PyResult<String> {
3357 self.ensure_registered_for_data()?;
3358 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3359 match params {
3360 Some(dict) => from_pydict(py, &dict),
3361 None => Ok(None),
3362 }
3363 })?;
3364 let request_id = DataActor::request_instruments(
3365 self.inner_mut(),
3366 venue,
3367 start,
3368 end,
3369 client_id,
3370 params_map,
3371 )
3372 .map_err(to_pyvalue_err)?;
3373 Ok(request_id.to_string())
3374 }
3375
3376 #[pyo3(name = "request_book_snapshot")]
3377 #[pyo3(signature = (instrument_id, depth=None, client_id=None, params=None))]
3378 fn py_request_book_snapshot(
3379 &mut self,
3380 instrument_id: InstrumentId,
3381 depth: Option<usize>,
3382 client_id: Option<ClientId>,
3383 params: Option<Py<PyDict>>,
3384 ) -> PyResult<String> {
3385 self.ensure_registered_for_data()?;
3386 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3387 match params {
3388 Some(dict) => from_pydict(py, &dict),
3389 None => Ok(None),
3390 }
3391 })?;
3392 let depth = depth.and_then(NonZeroUsize::new);
3393
3394 let request_id = DataActor::request_book_snapshot(
3395 self.inner_mut(),
3396 instrument_id,
3397 depth,
3398 client_id,
3399 params_map,
3400 )
3401 .map_err(to_pyvalue_err)?;
3402 Ok(request_id.to_string())
3403 }
3404
3405 #[pyo3(name = "request_book_deltas")]
3406 #[pyo3(signature = (instrument_id, start=None, end=None, limit=None, client_id=None, params=None))]
3407 fn py_request_book_deltas(
3408 &mut self,
3409 instrument_id: InstrumentId,
3410 start: Option<Timestamp>,
3411 end: Option<Timestamp>,
3412 limit: Option<usize>,
3413 client_id: Option<ClientId>,
3414 params: Option<Py<PyDict>>,
3415 ) -> PyResult<String> {
3416 self.ensure_registered_for_data()?;
3417 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3418 match params {
3419 Some(dict) => from_pydict(py, &dict),
3420 None => Ok(None),
3421 }
3422 })?;
3423 let limit = limit.and_then(NonZeroUsize::new);
3424 let request_id = DataActor::request_book_deltas(
3425 self.inner_mut(),
3426 instrument_id,
3427 start,
3428 end,
3429 limit,
3430 client_id,
3431 params_map,
3432 )
3433 .map_err(to_pyvalue_err)?;
3434 Ok(request_id.to_string())
3435 }
3436
3437 #[pyo3(name = "request_book_depth")]
3438 #[pyo3(signature = (instrument_id, start=None, end=None, limit=None, depth=None, client_id=None, params=None))]
3439 #[expect(clippy::too_many_arguments)]
3440 fn py_request_book_depth(
3441 &mut self,
3442 instrument_id: InstrumentId,
3443 start: Option<Timestamp>,
3444 end: Option<Timestamp>,
3445 limit: Option<usize>,
3446 depth: Option<usize>,
3447 client_id: Option<ClientId>,
3448 params: Option<Py<PyDict>>,
3449 ) -> PyResult<String> {
3450 self.ensure_registered_for_data()?;
3451 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3452 match params {
3453 Some(dict) => from_pydict(py, &dict),
3454 None => Ok(None),
3455 }
3456 })?;
3457 let limit = limit.and_then(NonZeroUsize::new);
3458 let depth = depth.and_then(NonZeroUsize::new);
3459 let request_id = DataActor::request_book_depth(
3460 self.inner_mut(),
3461 instrument_id,
3462 start,
3463 end,
3464 limit,
3465 depth,
3466 client_id,
3467 params_map,
3468 )
3469 .map_err(to_pyvalue_err)?;
3470 Ok(request_id.to_string())
3471 }
3472
3473 #[pyo3(name = "request_quotes")]
3474 #[pyo3(signature = (instrument_id, start=None, end=None, limit=None, client_id=None, params=None))]
3475 fn py_request_quotes(
3476 &mut self,
3477 instrument_id: InstrumentId,
3478 start: Option<Timestamp>,
3479 end: Option<Timestamp>,
3480 limit: Option<usize>,
3481 client_id: Option<ClientId>,
3482 params: Option<Py<PyDict>>,
3483 ) -> PyResult<String> {
3484 self.ensure_registered_for_data()?;
3485 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3486 match params {
3487 Some(dict) => from_pydict(py, &dict),
3488 None => Ok(None),
3489 }
3490 })?;
3491 let limit = limit.and_then(NonZeroUsize::new);
3492 let request_id = DataActor::request_quotes(
3493 self.inner_mut(),
3494 instrument_id,
3495 start,
3496 end,
3497 limit,
3498 client_id,
3499 params_map,
3500 )
3501 .map_err(to_pyvalue_err)?;
3502 Ok(request_id.to_string())
3503 }
3504
3505 #[pyo3(name = "request_trades")]
3506 #[pyo3(signature = (instrument_id, start=None, end=None, limit=None, client_id=None, params=None))]
3507 fn py_request_trades(
3508 &mut self,
3509 instrument_id: InstrumentId,
3510 start: Option<Timestamp>,
3511 end: Option<Timestamp>,
3512 limit: Option<usize>,
3513 client_id: Option<ClientId>,
3514 params: Option<Py<PyDict>>,
3515 ) -> PyResult<String> {
3516 self.ensure_registered_for_data()?;
3517 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3518 match params {
3519 Some(dict) => from_pydict(py, &dict),
3520 None => Ok(None),
3521 }
3522 })?;
3523 let limit = limit.and_then(NonZeroUsize::new);
3524 let request_id = DataActor::request_trades(
3525 self.inner_mut(),
3526 instrument_id,
3527 start,
3528 end,
3529 limit,
3530 client_id,
3531 params_map,
3532 )
3533 .map_err(to_pyvalue_err)?;
3534 Ok(request_id.to_string())
3535 }
3536
3537 #[pyo3(name = "request_funding_rates")]
3538 #[pyo3(signature = (instrument_id, start=None, end=None, limit=None, client_id=None, params=None))]
3539 fn py_request_funding_rates(
3540 &mut self,
3541 instrument_id: InstrumentId,
3542 start: Option<Timestamp>,
3543 end: Option<Timestamp>,
3544 limit: Option<usize>,
3545 client_id: Option<ClientId>,
3546 params: Option<Py<PyDict>>,
3547 ) -> PyResult<String> {
3548 self.ensure_registered_for_data()?;
3549 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3550 match params {
3551 Some(dict) => from_pydict(py, &dict),
3552 None => Ok(None),
3553 }
3554 })?;
3555 let limit = limit.and_then(NonZeroUsize::new);
3556 let request_id = DataActor::request_funding_rates(
3557 self.inner_mut(),
3558 instrument_id,
3559 start,
3560 end,
3561 limit,
3562 client_id,
3563 params_map,
3564 )
3565 .map_err(to_pyvalue_err)?;
3566 Ok(request_id.to_string())
3567 }
3568
3569 #[pyo3(name = "request_bars")]
3570 #[pyo3(signature = (bar_type, start=None, end=None, limit=None, client_id=None, params=None))]
3571 fn py_request_bars(
3572 &mut self,
3573 bar_type: BarType,
3574 start: Option<Timestamp>,
3575 end: Option<Timestamp>,
3576 limit: Option<usize>,
3577 client_id: Option<ClientId>,
3578 params: Option<Py<PyDict>>,
3579 ) -> PyResult<String> {
3580 self.ensure_registered_for_data()?;
3581 let params_map = Python::attach(|py| -> PyResult<Option<Params>> {
3582 match params {
3583 Some(dict) => from_pydict(py, &dict),
3584 None => Ok(None),
3585 }
3586 })?;
3587 let limit = limit.and_then(NonZeroUsize::new);
3588 let request_id = DataActor::request_bars(
3589 self.inner_mut(),
3590 bar_type,
3591 start,
3592 end,
3593 limit,
3594 client_id,
3595 params_map,
3596 )
3597 .map_err(to_pyvalue_err)?;
3598 Ok(request_id.to_string())
3599 }
3600
3601 #[pyo3(name = "reconnect_socket")]
3603 fn py_reconnect_socket(&self, client_id: ClientId, endpoint: &str) -> PyResult<()> {
3604 DataActor::reconnect_socket(self.inner(), client_id, endpoint).map_err(to_pyruntime_err)
3605 }
3606}
3607
3608#[pyo3_stub_gen::derive::gen_stub_pymethods]
3609#[pyo3::pymethods]
3610impl PyStrategy {
3611 #[pyo3(name = "publish_message", signature = (topic, message))]
3612 fn py_publish_message(
3613 slf: &Bound<'_, Self>,
3614 topic: &str,
3615 #[gen_stub(override_type(type_repr = "object"))] message: Py<PyAny>,
3616 ) -> PyResult<()> {
3617 let messages = get_python_message_bus(slf.as_any())?;
3618 messages.publish_message(topic, message)
3619 }
3620
3621 #[pyo3(name = "subscribe_topic")]
3622 #[pyo3(signature = (topic, handler, priority=0))]
3623 fn py_subscribe_topic(
3624 slf: &Bound<'_, Self>,
3625 topic: &str,
3626 #[gen_stub(override_type(type_repr = "collections.abc.Callable[[object], None]", imports = ("collections.abc",)))]
3627 handler: Py<PyAny>,
3628 priority: u32,
3629 ) -> PyResult<()> {
3630 let messages = get_python_message_bus(slf.as_any())?;
3631 messages.subscribe_topic(slf.py(), topic, handler, priority)
3632 }
3633
3634 #[pyo3(name = "unsubscribe_topic", signature = (topic, handler))]
3635 fn py_unsubscribe_topic(
3636 slf: &Bound<'_, Self>,
3637 topic: &str,
3638 #[gen_stub(override_type(type_repr = "collections.abc.Callable[[object], None]", imports = ("collections.abc",)))]
3639 handler: &Bound<'_, PyAny>,
3640 ) -> PyResult<()> {
3641 let messages = get_python_message_bus(slf.as_any())?;
3642 messages.unsubscribe_topic(topic, handler)
3643 }
3644}
3645
3646impl PyStrategy {
3647 fn ensure_registered_for_data(&self) -> PyResult<()> {
3648 if self.inner().core.actor.is_registered() {
3649 Ok(())
3650 } else {
3651 Err(to_pyruntime_err(
3652 "Strategy must be registered before publishing, managing synthetics, or requesting data",
3653 ))
3654 }
3655 }
3656
3657 fn ensure_registered(&self) -> PyResult<()> {
3658 if self.inner().core.actor.is_registered() {
3659 Ok(())
3660 } else {
3661 Err(to_pyruntime_err(
3662 "Strategy must be registered before managing subscriptions",
3663 ))
3664 }
3665 }
3666}
3667
3668fn class_derived_actor_id(slf: &Bound<'_, PyStrategy>, class_name: &str) -> PyResult<ActorId> {
3673 let borrowed = slf.borrow();
3674 let order_id_tag = normalize_order_id_tag(borrowed.inner().core.order_id_tag())
3675 .unwrap_or(UNASSIGNED_ORDER_ID_TAG);
3676
3677 ActorId::new_checked(format!("{class_name}-{order_id_tag}")).map_err(to_pyvalue_err)
3678}
3679
3680fn has_configured_strategy_id(slf: &Bound<'_, PyStrategy>) -> bool {
3686 let py = slf.py();
3687 let config = slf
3688 .borrow()
3689 .inner()
3690 .config
3691 .as_ref()
3692 .map(|config| config.clone_ref(py));
3693
3694 config.is_some_and(|config| {
3695 config
3696 .bind(py)
3697 .getattr("strategy_id")
3698 .is_ok_and(|strategy_id| !strategy_id.is_none())
3699 })
3700}
3701
3702fn py_order_list_to_orders(py: Python<'_>, order_list: &Py<PyAny>) -> PyResult<Vec<OrderAny>> {
3703 let order_objects = match order_list.getattr(py, "orders") {
3704 Ok(orders) => orders.extract::<Vec<Py<PyAny>>>(py)?,
3705 Err(e) if e.is_instance_of::<pyo3::exceptions::PyAttributeError>(py) => {
3706 order_list.extract::<Vec<Py<PyAny>>>(py)?
3707 }
3708 Err(e) => return Err(e),
3709 };
3710
3711 order_objects
3712 .into_iter()
3713 .map(|order| pyobject_to_order_any(py, order))
3714 .collect()
3715}
3716
3717#[cfg(test)]
3718mod tests {
3719 use std::{
3720 cell::RefCell,
3721 collections::{BTreeMap, HashMap},
3722 num::NonZeroUsize,
3723 rc::Rc,
3724 str::FromStr,
3725 };
3726
3727 use indexmap::IndexMap;
3728 use nautilus_common::{
3729 actor::{DataActor, DataActorNative, registry::actor_exists},
3730 cache::Cache,
3731 clock::{Clock, VirtualClock},
3732 component::{Component, get_component},
3733 enums::ComponentState,
3734 live::runner::replace_system_command_sender,
3735 logging::{
3736 arm_shutdown_on_error, disarm_shutdown_on_error, init_logging,
3737 take_shutdown_on_error_trigger,
3738 },
3739 messages::{
3740 SystemCommand,
3741 data::{
3742 BarsResponse, DataCommand, QuotesResponse, SubscribeCommand, TradesResponse,
3743 UnsubscribeCommand,
3744 },
3745 execution::TradingCommand,
3746 system::{
3747 QueueCondition, QueueState, QueueStateChanged, SocketState, SocketStateChanged,
3748 },
3749 },
3750 msgbus::{
3751 self, MessagingSwitchboard,
3752 stubs::{TypedIntoMessageSavingHandler, get_typed_into_message_saving_handler},
3753 },
3754 python::{cache::PyCache, wrappers::get_python_wrapper},
3755 runner::SystemChannel,
3756 signal::Signal,
3757 timer::TimeEvent,
3758 };
3759 use nautilus_core::{DurationNanos, UUID4, UnixNanos};
3760 use nautilus_model::{
3761 data::{
3762 Bar, BarType, CustomData, FundingRateUpdate, IndexPriceUpdate, InstrumentStatus,
3763 MarkPriceUpdate, OrderBookDelta, OrderBookDeltas, OrderBookDepth, QuoteTick, TradeTick,
3764 close::InstrumentClose,
3765 greeks::OptionGreekValues,
3766 option_chain::{OptionChainSlice, OptionGreeks},
3767 stubs::{stub_custom_data, stub_deltas, stub_depth10},
3768 },
3769 enums::{
3770 AggressorSide, BookType, GreeksConvention, InstrumentCloseType, MarketStatusAction,
3771 OmsType, OrderSide, OrderType, PositionSide, TimeInForce,
3772 },
3773 events::{
3774 OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
3775 OrderEventAny, OrderExpired, OrderInitialized, OrderModifyRejected, OrderPendingCancel,
3776 OrderPendingUpdate, OrderRejected, OrderReleased, OrderSubmitted, OrderTriggered,
3777 OrderUpdated, PositionChanged, PositionClosed, PositionEvent, PositionOpened,
3778 order::spec::{OrderFillVoidedSpec, OrderFilledSpec},
3779 },
3780 identifiers::{
3781 AccountId, ClientId, ClientOrderId, ComponentId, InstrumentId, OptionSeriesId,
3782 OrderListId, PositionId, StrategyId, TradeId, TraderId, Venue,
3783 },
3784 instruments::{CurrencyPair, InstrumentAny, stubs::audusd_sim},
3785 orderbook::OrderBook,
3786 orders::{Order, OrderTestBuilder},
3787 position::Position,
3788 python::orders::order_any_to_pyobject,
3789 types::{Currency, Money, Price, Quantity},
3790 };
3791 use nautilus_portfolio::portfolio::Portfolio;
3792 use pyo3::{
3793 Bound, Py, PyAny, PyResult, Python,
3794 ffi::c_str,
3795 types::{PyAnyMethods, PyBytes, PyDict, PyList, PyWeakrefMethods, PyWeakrefReference},
3796 };
3797 use serde_json::Value;
3798 use ustr::Ustr;
3799
3800 use super::PyStrategy;
3801 use crate::strategy::{Strategy, StrategyConfig};
3802
3803 const TRACKING_STRATEGY_CODE: &std::ffi::CStr = c_str!(
3804 r#"
3805class TrackingStrategy:
3806 TRACKED_METHODS = {
3807 "on_start",
3808 "on_stop",
3809 "on_resume",
3810 "on_reset",
3811 "on_dispose",
3812 "on_degrade",
3813 "on_fault",
3814 "on_save",
3815 "on_load",
3816 "on_time_event",
3817 "on_data",
3818 "on_signal",
3819 "on_queue_state",
3820 "on_socket_state",
3821 "on_instrument",
3822 "on_quote",
3823 "on_trade",
3824 "on_bar",
3825 "on_book_deltas",
3826 "on_book_depth",
3827 "on_book",
3828 "on_mark_price",
3829 "on_index_price",
3830 "on_funding_rate",
3831 "on_instrument_status",
3832 "on_instrument_close",
3833 "on_option_greeks",
3834 "on_option_chain",
3835 "on_historical_data",
3836 "on_historical_book_deltas",
3837 "on_historical_book_depth",
3838 "on_historical_quotes",
3839 "on_historical_trades",
3840 "on_historical_funding_rates",
3841 "on_historical_bars",
3842 "on_historical_mark_prices",
3843 "on_historical_index_prices",
3844 "on_market_exit",
3845 "post_market_exit",
3846 "on_order_initialized",
3847 "on_order_event",
3848 "on_order_denied",
3849 "on_order_emulated",
3850 "on_order_released",
3851 "on_order_submitted",
3852 "on_order_rejected",
3853 "on_order_accepted",
3854 "on_order_expired",
3855 "on_order_triggered",
3856 "on_order_pending_update",
3857 "on_order_pending_cancel",
3858 "on_order_modify_rejected",
3859 "on_order_cancel_rejected",
3860 "on_order_updated",
3861 "on_order_canceled",
3862 "on_order_filled",
3863 "on_order_fill_voided",
3864 "on_position_opened",
3865 "on_position_event",
3866 "on_position_changed",
3867 "on_position_closed",
3868 }
3869
3870 def __init__(self):
3871 self.calls = []
3872 self.raise_on = None
3873
3874 def _record(self, method_name, *args):
3875 self.calls.append((method_name, args))
3876 if self.raise_on == method_name:
3877 self.fail_callback()
3878
3879 def fail_callback(self):
3880 try:
3881 raise ValueError("callback cause")
3882 except ValueError as e:
3883 raise RuntimeError("callback failure") from e
3884
3885 def was_called(self, method_name):
3886 return any(call[0] == method_name for call in self.calls)
3887
3888 def call_count(self, method_name):
3889 return sum(1 for call in self.calls if call[0] == method_name)
3890
3891 def last_call_args(self, method_name):
3892 for called_method, args in reversed(self.calls):
3893 if called_method == method_name:
3894 return args
3895 raise AssertionError(f"{method_name} was not called")
3896
3897 def call_names(self):
3898 return [call[0] for call in self.calls]
3899
3900 def last_loaded_state(self):
3901
3902 for method_name, args in reversed(self.calls):
3903 if method_name == "on_load":
3904 return args[0]
3905 return None
3906
3907 def on_save(self):
3908 self._record("on_save")
3909 return {"strategy": b"saved"}
3910
3911 def on_load(self, state):
3912 self._record("on_load", dict(state))
3913
3914 def __getattr__(self, name):
3915 if name in self.TRACKED_METHODS:
3916 return lambda *args: self._record(name, *args)
3917 raise AttributeError(name)
3918"#
3919 );
3920
3921 fn create_tracking_python_strategy(py: Python<'_>) -> PyResult<Py<PyAny>> {
3922 py.run(TRACKING_STRATEGY_CODE, None, None)?;
3923 let tracking_strategy_class = py.eval(c_str!("TrackingStrategy"), None, None)?;
3924 let instance = tracking_strategy_class.call0()?;
3925 Ok(instance.unbind())
3926 }
3927
3928 fn python_method_was_called(
3929 py_strategy: &Py<PyAny>,
3930 py: Python<'_>,
3931 method_name: &str,
3932 ) -> bool {
3933 py_strategy
3934 .call_method1(py, "was_called", (method_name,))
3935 .and_then(|result| result.extract::<bool>(py))
3936 .unwrap_or(false)
3937 }
3938
3939 fn python_method_call_count(py_strategy: &Py<PyAny>, py: Python<'_>, method_name: &str) -> i32 {
3940 py_strategy
3941 .call_method1(py, "call_count", (method_name,))
3942 .and_then(|result| result.extract::<i32>(py))
3943 .unwrap_or(0)
3944 }
3945
3946 fn python_method_call_names(py_strategy: &Py<PyAny>, py: Python<'_>) -> Vec<String> {
3947 py_strategy
3948 .call_method0(py, "call_names")
3949 .and_then(|result| result.extract::<Vec<String>>(py))
3950 .unwrap_or_default()
3951 }
3952
3953 fn python_last_loaded_state(
3954 py_strategy: &Py<PyAny>,
3955 py: Python<'_>,
3956 ) -> Option<HashMap<String, Vec<u8>>> {
3957 py_strategy
3958 .call_method0(py, "last_loaded_state")
3959 .and_then(|result| result.extract::<Option<HashMap<String, Vec<u8>>>>(py))
3960 .unwrap_or(None)
3961 }
3962
3963 const TRACKING_INDICATOR_CODE: &std::ffi::CStr = c_str!(
3964 r#"
3965class TrackingIndicator:
3966 def __init__(self, events=None):
3967 self.initialized = False
3968 self.calls = []
3969 self.events = events
3970
3971 def handle_quote_tick(self, quote):
3972 self.calls.append("quote")
3973 if self.events is not None:
3974 self.events.append("indicator:quote")
3975
3976 def handle_trade_tick(self, trade):
3977 self.calls.append("trade")
3978 if self.events is not None:
3979 self.events.append("indicator:trade")
3980
3981 def handle_bar(self, bar):
3982 self.calls.append("bar")
3983 if self.events is not None:
3984 self.events.append("indicator:bar")
3985
3986 def call_count(self, name):
3987 return self.calls.count(name)
3988
3989class IndicatorEventStrategy:
3990 def __init__(self, events):
3991 self.events = events
3992
3993 def on_start(self):
3994 pass
3995
3996 def on_quote(self, quote):
3997 self.events.append("strategy:quote")
3998
3999 def on_trade(self, trade):
4000 self.events.append("strategy:trade")
4001
4002 def on_bar(self, bar):
4003 self.events.append("strategy:bar")
4004"#
4005 );
4006
4007 fn create_tracking_python_indicator(py: Python<'_>) -> PyResult<Py<PyAny>> {
4008 py.run(TRACKING_INDICATOR_CODE, None, None)?;
4009 let indicator_class = py.eval(c_str!("TrackingIndicator"), None, None)?;
4010 Ok(indicator_class.call0()?.unbind())
4011 }
4012
4013 fn create_event_tracking_python_indicator(
4014 py: Python<'_>,
4015 events: &Bound<'_, PyList>,
4016 ) -> PyResult<Py<PyAny>> {
4017 py.run(TRACKING_INDICATOR_CODE, None, None)?;
4018 let indicator_class = py.eval(c_str!("TrackingIndicator"), None, None)?;
4019 Ok(indicator_class.call1((events,))?.unbind())
4020 }
4021
4022 fn create_indicator_event_strategy(
4023 py: Python<'_>,
4024 events: &Bound<'_, PyList>,
4025 ) -> PyResult<Py<PyAny>> {
4026 py.run(TRACKING_INDICATOR_CODE, None, None)?;
4027 let strategy_class = py.eval(c_str!("IndicatorEventStrategy"), None, None)?;
4028 Ok(strategy_class.call1((events,))?.unbind())
4029 }
4030
4031 fn python_indicator_call_count(
4032 indicator: &Py<PyAny>,
4033 py: Python<'_>,
4034 method_name: &str,
4035 ) -> i32 {
4036 indicator
4037 .call_method1(py, "call_count", (method_name,))
4038 .and_then(|result| result.extract::<i32>(py))
4039 .unwrap_or(0)
4040 }
4041
4042 fn sample_instrument() -> CurrencyPair {
4043 audusd_sim()
4044 }
4045
4046 fn sample_time_event() -> TimeEvent {
4047 TimeEvent::new(
4048 Ustr::from("test_timer"),
4049 UUID4::new(),
4050 UnixNanos::default(),
4051 UnixNanos::default(),
4052 )
4053 }
4054
4055 fn sample_data() -> CustomData {
4056 stub_custom_data(1, 42, None, None)
4057 }
4058
4059 fn sample_signal() -> Signal {
4060 Signal::new(
4061 Ustr::from("test_signal"),
4062 "1.0".to_string(),
4063 UnixNanos::default(),
4064 UnixNanos::default(),
4065 )
4066 }
4067
4068 fn sample_queue_state_changed() -> QueueStateChanged {
4069 QueueStateChanged::new(
4070 TraderId::from("TRADER-001"),
4071 SystemChannel::ExecCommands,
4072 QueueCondition::Backlogged,
4073 QueueState::Triggered,
4074 17,
4075 23,
4076 UUID4::from("00000000-0000-4000-8000-000000000001"),
4077 UnixNanos::from(1_700_000_000_000_000_001),
4078 UnixNanos::from(1_700_000_000_000_000_002),
4079 )
4080 }
4081
4082 fn sample_socket_state_changed() -> SocketStateChanged {
4083 SocketStateChanged::new(
4084 TraderId::from("TRADER-001"),
4085 ClientId::from("BINANCE"),
4086 Some(Venue::from("BINANCE")),
4087 Ustr::from("binance-futures-market-streams"),
4088 SocketState::Connected,
4089 UUID4::from("00000000-0000-4000-8000-000000000001"),
4090 UnixNanos::from(1_700_000_000_000_000_001),
4091 UnixNanos::from(1_700_000_000_000_000_002),
4092 )
4093 }
4094
4095 fn sample_quote() -> QuoteTick {
4096 let instrument = sample_instrument();
4097 QuoteTick::new(
4098 instrument.id,
4099 Price::from("1.00000"),
4100 Price::from("1.00001"),
4101 Quantity::from(100_000),
4102 Quantity::from(100_000),
4103 UnixNanos::default(),
4104 UnixNanos::default(),
4105 )
4106 }
4107
4108 fn sample_trade() -> TradeTick {
4109 let instrument = sample_instrument();
4110 TradeTick::new(
4111 instrument.id,
4112 Price::from("1.00000"),
4113 Quantity::from(100_000),
4114 AggressorSide::Buy,
4115 TradeId::new("123456"),
4116 UnixNanos::default(),
4117 UnixNanos::default(),
4118 )
4119 }
4120
4121 fn sample_bar() -> Bar {
4122 let instrument = sample_instrument();
4123 let bar_type =
4124 BarType::from_str(&format!("{}-1-MINUTE-LAST-INTERNAL", instrument.id)).unwrap();
4125 Bar::new(
4126 bar_type,
4127 Price::from("1.00000"),
4128 Price::from("1.00010"),
4129 Price::from("0.99990"),
4130 Price::from("1.00005"),
4131 Quantity::from(100_000),
4132 UnixNanos::default(),
4133 UnixNanos::default(),
4134 )
4135 }
4136
4137 fn sample_book() -> OrderBook {
4138 OrderBook::new(sample_instrument().id, BookType::L2_MBP)
4139 }
4140
4141 fn sample_book_deltas() -> OrderBookDeltas {
4142 let instrument = sample_instrument();
4143 let delta =
4144 OrderBookDelta::clear(instrument.id, 0, UnixNanos::default(), UnixNanos::default());
4145 OrderBookDeltas::new(instrument.id, vec![delta])
4146 }
4147
4148 fn sample_book_depth() -> OrderBookDepth {
4149 stub_depth10()
4150 }
4151
4152 fn sample_mark_price() -> MarkPriceUpdate {
4153 MarkPriceUpdate::new(
4154 sample_instrument().id,
4155 Price::from("1.00000"),
4156 UnixNanos::default(),
4157 UnixNanos::default(),
4158 )
4159 }
4160
4161 fn sample_index_price() -> IndexPriceUpdate {
4162 IndexPriceUpdate::new(
4163 sample_instrument().id,
4164 Price::from("1.00000"),
4165 UnixNanos::default(),
4166 UnixNanos::default(),
4167 )
4168 }
4169
4170 fn sample_funding_rate() -> FundingRateUpdate {
4171 FundingRateUpdate::new(
4172 sample_instrument().id,
4173 "0.0001".parse().unwrap(),
4174 None,
4175 None,
4176 UnixNanos::default(),
4177 UnixNanos::default(),
4178 )
4179 }
4180
4181 fn sample_instrument_status() -> InstrumentStatus {
4182 InstrumentStatus::new(
4183 sample_instrument().id,
4184 MarketStatusAction::Trading,
4185 UnixNanos::default(),
4186 UnixNanos::default(),
4187 None,
4188 None,
4189 None,
4190 None,
4191 None,
4192 )
4193 }
4194
4195 fn sample_instrument_close() -> InstrumentClose {
4196 InstrumentClose::new(
4197 sample_instrument().id,
4198 Price::from("1.00000"),
4199 InstrumentCloseType::EndOfSession,
4200 UnixNanos::default(),
4201 UnixNanos::default(),
4202 )
4203 }
4204
4205 fn sample_option_greeks() -> OptionGreeks {
4206 OptionGreeks {
4207 instrument_id: InstrumentId::from("AUD/USD.SIM"),
4208 convention: GreeksConvention::BlackScholes,
4209 greeks: OptionGreekValues {
4210 delta: 0.55,
4211 gamma: 0.03,
4212 vega: 0.12,
4213 theta: -0.05,
4214 rho: 0.01,
4215 },
4216 mark_iv: Some(0.25),
4217 bid_iv: None,
4218 ask_iv: None,
4219 underlying_price: None,
4220 open_interest: None,
4221 ts_event: UnixNanos::default(),
4222 ts_init: UnixNanos::default(),
4223 }
4224 }
4225
4226 fn sample_option_chain() -> OptionChainSlice {
4227 OptionChainSlice {
4228 series_id: OptionSeriesId::new(
4229 Venue::from("SIM"),
4230 Ustr::from("AUD"),
4231 Ustr::from("USD"),
4232 UnixNanos::from(1_711_036_800_000_000_000),
4233 ),
4234 atm_strike: None,
4235 calls: BTreeMap::default(),
4236 puts: BTreeMap::default(),
4237 ts_event: UnixNanos::default(),
4238 ts_init: UnixNanos::default(),
4239 }
4240 }
4241
4242 fn sample_position_opened() -> PositionOpened {
4243 PositionOpened {
4244 trader_id: TraderId::from("TRADER-001"),
4245 strategy_id: StrategyId::from("TEST-001"),
4246 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
4247 position_id: PositionId::from("P-001"),
4248 account_id: AccountId::from("ACC-001"),
4249 opening_order_id: ClientOrderId::from("O-001"),
4250 entry: OrderSide::Buy,
4251 side: PositionSide::Long,
4252 signed_qty: 1.0,
4253 quantity: Quantity::from(1),
4254 last_qty: Quantity::from(1),
4255 last_px: Price::from("1.00000"),
4256 currency: Currency::from("USD"),
4257 avg_px_open: 1.0,
4258 realized_pnl: None,
4259 event_id: UUID4::new(),
4260 ts_event: UnixNanos::default(),
4261 ts_init: UnixNanos::default(),
4262 }
4263 }
4264
4265 fn sample_position_changed() -> PositionChanged {
4266 PositionChanged {
4267 trader_id: TraderId::from("TRADER-001"),
4268 strategy_id: StrategyId::from("TEST-001"),
4269 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
4270 position_id: PositionId::from("P-001"),
4271 account_id: AccountId::from("ACC-001"),
4272 opening_order_id: ClientOrderId::from("O-001"),
4273 entry: OrderSide::Buy,
4274 side: PositionSide::Long,
4275 signed_qty: 2.0,
4276 quantity: Quantity::from(2),
4277 peak_quantity: Quantity::from(2),
4278 last_qty: Quantity::from(1),
4279 last_px: Price::from("1.10000"),
4280 currency: Currency::from("USD"),
4281 avg_px_open: 1.05,
4282 avg_px_close: None,
4283 realized_return: 0.0,
4284 realized_pnl: None,
4285 unrealized_pnl: Money::new(0.0, Currency::USD()),
4286 event_id: UUID4::new(),
4287 ts_opened: UnixNanos::default(),
4288 ts_event: UnixNanos::default(),
4289 ts_init: UnixNanos::default(),
4290 }
4291 }
4292
4293 fn sample_position_closed() -> PositionClosed {
4294 PositionClosed {
4295 trader_id: TraderId::from("TRADER-001"),
4296 strategy_id: StrategyId::from("TEST-001"),
4297 instrument_id: InstrumentId::from("BTCUSDT.BINANCE"),
4298 position_id: PositionId::from("P-001"),
4299 account_id: AccountId::from("ACC-001"),
4300 opening_order_id: ClientOrderId::from("O-001"),
4301 closing_order_id: Some(ClientOrderId::from("O-002")),
4302 entry: OrderSide::Buy,
4303 side: PositionSide::Flat,
4304 signed_qty: 0.0,
4305 quantity: Quantity::from(0),
4306 peak_quantity: Quantity::from(2),
4307 last_qty: Quantity::from(2),
4308 last_px: Price::from("1.20000"),
4309 currency: Currency::from("USD"),
4310 avg_px_open: 1.05,
4311 avg_px_close: Some(1.20),
4312 realized_return: 0.1,
4313 realized_pnl: Some(Money::new(0.1, Currency::USD())),
4314 unrealized_pnl: Money::new(0.0, Currency::USD()),
4315 duration: DurationNanos::new(1),
4316 event_id: UUID4::new(),
4317 ts_opened: UnixNanos::default(),
4318 ts_closed: Some(UnixNanos::default()),
4319 ts_event: UnixNanos::default(),
4320 ts_init: UnixNanos::default(),
4321 }
4322 }
4323
4324 fn sample_python_market_order(
4325 py: Python<'_>,
4326 strategy_id: StrategyId,
4327 client_order_id: ClientOrderId,
4328 ) -> PyResult<Py<PyAny>> {
4329 let order = OrderTestBuilder::new(OrderType::Market)
4330 .trader_id(TraderId::from("TRADER-001"))
4331 .strategy_id(strategy_id)
4332 .instrument_id(sample_instrument().id)
4333 .client_order_id(client_order_id)
4334 .quantity(Quantity::from(100_000))
4335 .build();
4336
4337 order_any_to_pyobject(py, order)
4338 }
4339
4340 fn sample_open_position(
4341 strategy_id: StrategyId,
4342 position_id: PositionId,
4343 client_order_id: ClientOrderId,
4344 ) -> Position {
4345 let instrument = sample_instrument();
4346 let fill = OrderFilledSpec::builder()
4347 .trader_id(TraderId::from("TRADER-001"))
4348 .strategy_id(strategy_id)
4349 .instrument_id(instrument.id)
4350 .client_order_id(client_order_id)
4351 .position_id(position_id)
4352 .order_side(OrderSide::Buy)
4353 .last_qty(Quantity::from(100_000))
4354 .last_px(Price::from("1.00000"))
4355 .currency(instrument.quote_currency)
4356 .build();
4357
4358 Position::new(&InstrumentAny::CurrencyPair(instrument), fill)
4359 }
4360
4361 fn create_registered_tracking_strategy_with_config(
4362 py: Python<'_>,
4363 config: Option<StrategyConfig>,
4364 ) -> (Py<PyAny>, PyStrategy) {
4365 let py_strategy = create_tracking_python_strategy(py).unwrap();
4366 let mut rust_strategy = PyStrategy::new(config);
4367 rust_strategy
4368 .set_python_instance(py_strategy.bind(py))
4369 .unwrap();
4370
4371 let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
4372 let cache = Rc::new(RefCell::new(Cache::new(None, None)));
4373 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4374 clock.clone(),
4375 cache.clone(),
4376 None,
4377 )));
4378
4379 rust_strategy
4380 .register(TraderId::from("TRADER-001"), clock, cache, portfolio)
4381 .unwrap();
4382
4383 (py_strategy, rust_strategy)
4384 }
4385
4386 fn create_registered_tracking_strategy(py: Python<'_>) -> (Py<PyAny>, PyStrategy) {
4387 create_registered_tracking_strategy_with_config(py, None)
4388 }
4389
4390 #[rstest::rstest]
4391 fn test_register_in_global_registries_retains_python_wrapper() {
4392 pyo3::Python::initialize();
4393
4394 Python::attach(|py| {
4395 let (py_strategy, rust_strategy) = create_registered_tracking_strategy_with_config(
4396 py,
4397 Some(StrategyConfig {
4398 strategy_id: Some(StrategyId::from("Retained-001")),
4399 ..Default::default()
4400 }),
4401 );
4402
4403 rust_strategy.register_in_global_registries().unwrap();
4404
4405 let retained = get_python_wrapper(ComponentId::from("Retained-001"))
4406 .expect("registering must retain the strategy's Python wrapper");
4407
4408 assert!(retained.bind(py).is(py_strategy.bind(py)));
4409 assert!(get_component(&Ustr::from("Retained-001")).is_some());
4410 assert!(actor_exists(&Ustr::from("Retained-001")));
4411 });
4412 }
4413
4414 #[rstest::rstest]
4415 fn test_registered_python_strategy_routes_time_events_by_state() {
4416 pyo3::Python::initialize();
4417
4418 Python::attach(|py| {
4419 let (py_strategy, rust_strategy) = create_registered_tracking_strategy_with_config(
4420 py,
4421 Some(StrategyConfig {
4422 strategy_id: Some(StrategyId::from("StoppedTimer-001")),
4423 manage_gtd_expiry: true,
4424 ..Default::default()
4425 }),
4426 );
4427 rust_strategy.register_in_global_registries().unwrap();
4428 Component::start(rust_strategy.inner_mut()).unwrap();
4429
4430 let running_client_order_id = ClientOrderId::from("O-RUNNING-001");
4431 let running_timer_name = format!("GTD-EXPIRY:{running_client_order_id}");
4432 let clock = DataActorNative::clock_rc(rust_strategy.inner());
4433 clock
4434 .borrow_mut()
4435 .set_time_alert_ns(&running_timer_name, UnixNanos::from(1), None, None)
4436 .unwrap();
4437 rust_strategy
4438 .inner_mut()
4439 .core
4440 .gtd_timers
4441 .insert(running_client_order_id, Ustr::from(&running_timer_name));
4442 let running_dispatched = dispatch_time_events(&clock, UnixNanos::from(1));
4443
4444 assert_eq!(running_dispatched, 1);
4445 assert!(
4446 !rust_strategy
4447 .inner_mut()
4448 .has_gtd_expiry_timer(&running_client_order_id)
4449 );
4450 assert_eq!(
4451 python_method_call_count(&py_strategy, py, "on_time_event"),
4452 1
4453 );
4454
4455 let client_order_id = ClientOrderId::from("O-STOPPED-001");
4456 let timer_name = format!("GTD-EXPIRY:{client_order_id}");
4457 clock
4458 .borrow_mut()
4459 .set_time_alert_ns(&timer_name, UnixNanos::from(2), None, None)
4460 .unwrap();
4461 rust_strategy
4462 .inner_mut()
4463 .core
4464 .gtd_timers
4465 .insert(client_order_id, Ustr::from(&timer_name));
4466 Component::stop(rust_strategy.inner_mut()).unwrap();
4467
4468 let stopped_dispatched = dispatch_time_events(&clock, UnixNanos::from(2));
4469
4470 assert_eq!(stopped_dispatched, 1);
4471 assert_eq!(
4472 Component::state(rust_strategy.inner()),
4473 ComponentState::Stopped
4474 );
4475 assert!(
4476 rust_strategy
4477 .inner_mut()
4478 .has_gtd_expiry_timer(&client_order_id)
4479 );
4480 assert_eq!(
4481 python_method_call_count(&py_strategy, py, "on_time_event"),
4482 1
4483 );
4484 });
4485 }
4486
4487 fn dispatch_time_events(clock: &Rc<RefCell<dyn Clock>>, to_time_ns: UnixNanos) -> usize {
4488 let handlers = {
4489 let mut clock_ref = clock.borrow_mut();
4490 let test_clock = clock_ref
4491 .as_any_mut()
4492 .downcast_mut::<VirtualClock>()
4493 .expect("strategy clock must be VirtualClock");
4494 let events = test_clock.advance_time(to_time_ns, true);
4495 test_clock.match_handlers(events)
4496 };
4497 let dispatched = handlers.len();
4498
4499 for handler in handlers {
4500 handler.run();
4501 }
4502
4503 dispatched
4504 }
4505
4506 #[rstest::rstest]
4507 fn test_register_in_global_registries_rejects_missing_python_wrapper() {
4508 pyo3::Python::initialize();
4509
4510 Python::attach(|_py| {
4511 let mut rust_strategy = PyStrategy::new(Some(StrategyConfig {
4512 strategy_id: Some(StrategyId::from("Unwrapped-001")),
4513 ..Default::default()
4514 }));
4515
4516 let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
4517 let cache = Rc::new(RefCell::new(Cache::new(None, None)));
4518 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4519 clock.clone(),
4520 cache.clone(),
4521 None,
4522 )));
4523 rust_strategy
4524 .register(TraderId::from("TRADER-001"), clock, cache, portfolio)
4525 .unwrap();
4526
4527 let error = rust_strategy
4528 .register_in_global_registries()
4529 .expect_err("registering without a Python wrapper must fail");
4530
4531 assert!(error.to_string().contains("without a Python wrapper"));
4532 assert!(get_component(&Ustr::from("Unwrapped-001")).is_none());
4533 assert!(!actor_exists(&Ustr::from("Unwrapped-001")));
4534 assert!(get_python_wrapper(ComponentId::from("Unwrapped-001")).is_none());
4535 });
4536 }
4537
4538 #[rstest::rstest]
4539 fn test_external_order_instrument_ids_returns_configured_instruments() {
4540 let claims = vec![
4541 InstrumentId::from("AUDUSD.SIM"),
4542 InstrumentId::from("BTCUSDT.BINANCE"),
4543 ];
4544 let strategy = PyStrategy::new(Some(StrategyConfig {
4545 external_order_instrument_ids: Some(claims.clone()),
4546 ..Default::default()
4547 }));
4548
4549 assert_eq!(strategy.external_order_instrument_ids(), Some(claims));
4550 }
4551
4552 #[rstest::rstest]
4553 fn test_new_checked_rejects_order_id_tag_with_separator() {
4554 let config = StrategyConfig {
4555 order_id_tag: Some("A-B".to_string()),
4556 ..Default::default()
4557 };
4558
4559 let error = PyStrategy::new_checked(Some(config)).unwrap_err();
4560
4561 assert_eq!(
4562 error.to_string(),
4563 "`order_id_tag` cannot contain the '-' strategy ID separator, was 'A-B'"
4564 );
4565 }
4566
4567 #[rstest::rstest]
4568 fn test_py_new_raises_value_error_for_order_id_tag_with_separator() {
4569 pyo3::Python::initialize();
4570 Python::attach(|py| {
4571 let config = StrategyConfig {
4574 order_id_tag: Some("A-B".to_string()),
4575 ..Default::default()
4576 };
4577 let config_obj = Py::new(py, config).unwrap().into_any();
4578
4579 let error = PyStrategy::py_new(Some(config_obj)).unwrap_err();
4580
4581 assert!(error.is_instance_of::<pyo3::exceptions::PyValueError>(py));
4582 assert_eq!(
4583 error.value(py).to_string(),
4584 "`order_id_tag` cannot contain the '-' strategy ID separator, was 'A-B'"
4585 );
4586 });
4587 }
4588
4589 #[rstest::rstest]
4590 fn test_py_new_accepts_a_config_without_a_separator() {
4591 pyo3::Python::initialize();
4592 Python::attach(|py| {
4593 let config = StrategyConfig {
4594 order_id_tag: Some("001".to_string()),
4595 ..Default::default()
4596 };
4597 let config_obj = Py::new(py, config).unwrap().into_any();
4598
4599 let strategy = PyStrategy::py_new(Some(config_obj)).unwrap();
4600
4601 assert_eq!(strategy.order_id_tag(), Some("001".to_string()));
4602 });
4603 }
4604
4605 #[rstest::rstest]
4606 fn test_python_aggregate_event_handlers_exist() {
4607 pyo3::Python::initialize();
4608 Python::attach(|py| {
4609 let strategy = Py::new(py, PyStrategy::new(None)).unwrap();
4610 let strategy = strategy.bind(py);
4611
4612 assert!(strategy.hasattr("on_order_event").unwrap());
4613 assert!(strategy.hasattr("on_position_event").unwrap());
4614 });
4615 }
4616
4617 #[rstest::rstest]
4618 #[case(None)]
4619 #[case(Some(25))]
4620 fn test_python_book_depth_subscription_methods_send_commands(#[case] depth: Option<usize>) {
4621 pyo3::Python::initialize();
4622 Python::attach(|py| {
4623 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
4624 let (handler, saver) = get_typed_into_message_saving_handler::<DataCommand>(None);
4625 msgbus::register_data_command_endpoint(
4626 MessagingSwitchboard::data_engine_queue_execute(),
4627 handler,
4628 );
4629
4630 let instrument_id = sample_instrument().id;
4631 let client_id = Some(ClientId::new("DEPTH-CLIENT"));
4632 rust_strategy
4633 .py_subscribe_book_depth(
4634 instrument_id,
4635 BookType::L2_MBP,
4636 depth,
4637 client_id,
4638 true,
4639 None,
4640 )
4641 .unwrap();
4642 rust_strategy
4643 .py_unsubscribe_book_depth(instrument_id, client_id, None)
4644 .unwrap();
4645
4646 let commands = saver.get_messages();
4647 let [
4648 DataCommand::Subscribe(SubscribeCommand::BookDepth(subscribe)),
4649 DataCommand::Unsubscribe(UnsubscribeCommand::BookDepth(unsubscribe)),
4650 ] = commands.as_slice()
4651 else {
4652 panic!("expected BookDepth subscribe and unsubscribe commands, was {commands:?}");
4653 };
4654
4655 assert_eq!(subscribe.instrument_id, instrument_id);
4656 assert_eq!(subscribe.book_type, BookType::L2_MBP);
4657 assert_eq!(subscribe.depth.map(NonZeroUsize::get), depth);
4658 assert_eq!(subscribe.client_id, client_id);
4659 assert!(subscribe.managed);
4660 assert_eq!(unsubscribe.instrument_id, instrument_id);
4661 assert_eq!(unsubscribe.client_id, client_id);
4662 });
4663 }
4664
4665 #[rstest::rstest]
4666 fn test_strategy_retains_python_config_object() {
4667 pyo3::Python::initialize();
4668 Python::attach(|py| {
4669 let config = py
4670 .eval(
4671 c_str!("type('_Cfg', (), {'strategy_id': 'S-RETAIN-001'})()"),
4672 None,
4673 None,
4674 )
4675 .unwrap();
4676
4677 let strategy = py
4678 .get_type::<PyStrategy>()
4679 .as_any()
4680 .call1((config.clone(),))
4681 .unwrap();
4682
4683 let retained = strategy.getattr("config").unwrap();
4684
4685 assert!(retained.is(&config));
4686 });
4687 }
4688
4689 #[rstest::rstest]
4690 fn test_indicator_registration_exposes_readiness_and_registered_view() {
4691 pyo3::Python::initialize();
4692 Python::attach(|py| {
4693 let mut rust_strategy = PyStrategy::new(None);
4694 let indicator = create_tracking_python_indicator(py).unwrap();
4695 let instrument_id = sample_instrument().id;
4696 let bar_type = sample_bar().bar_type;
4697
4698 assert_eq!(
4699 rust_strategy
4700 .py_registered_indicators(py)
4701 .unwrap()
4702 .bind(py)
4703 .len()
4704 .unwrap(),
4705 0
4706 );
4707 assert!(!rust_strategy.py_indicators_initialized(py).unwrap());
4708
4709 rust_strategy.py_register_indicator_for_quote_ticks(
4710 py,
4711 instrument_id,
4712 indicator.clone_ref(py),
4713 );
4714 rust_strategy.py_register_indicator_for_trade_ticks(
4715 py,
4716 instrument_id,
4717 indicator.clone_ref(py),
4718 );
4719 rust_strategy.py_register_indicator_for_bars(py, bar_type, indicator.clone_ref(py));
4720
4721 let registered = rust_strategy.py_registered_indicators(py).unwrap();
4722 let registered = registered.bind(py);
4723
4724 assert_eq!(registered.len().unwrap(), 1);
4725 assert_eq!(
4726 registered.get_item(0).unwrap().as_ptr(),
4727 indicator.bind(py).as_ptr()
4728 );
4729 assert!(!rust_strategy.py_indicators_initialized(py).unwrap());
4730
4731 indicator.bind(py).setattr("initialized", true).unwrap();
4732
4733 assert!(rust_strategy.py_indicators_initialized(py).unwrap());
4734 });
4735 }
4736
4737 #[rstest::rstest]
4738 fn test_registered_indicators_receive_quote_trade_and_bar_before_strategy_callbacks() {
4739 pyo3::Python::initialize();
4740 Python::attach(|py| {
4741 let events = PyList::empty(py);
4742 let py_strategy = create_indicator_event_strategy(py, &events).unwrap();
4743 let indicator = create_event_tracking_python_indicator(py, &events).unwrap();
4744
4745 let mut rust_strategy = PyStrategy::new(None);
4746 rust_strategy
4747 .set_python_instance(py_strategy.bind(py))
4748 .unwrap();
4749
4750 let clock: Rc<RefCell<dyn Clock>> = Rc::new(RefCell::new(VirtualClock::new()));
4751 let cache = Rc::new(RefCell::new(Cache::new(None, None)));
4752 let portfolio = Rc::new(RefCell::new(Portfolio::new(
4753 clock.clone(),
4754 cache.clone(),
4755 None,
4756 )));
4757
4758 rust_strategy
4759 .register(TraderId::from("TRADER-001"), clock, cache, portfolio)
4760 .unwrap();
4761 Component::start(rust_strategy.inner_mut()).unwrap();
4762
4763 let quote = sample_quote();
4764 let trade = sample_trade();
4765 let bar = sample_bar();
4766 let external_bar_type = BarType::from_str(&format!(
4767 "{}-1-MINUTE-LAST-EXTERNAL",
4768 bar.bar_type.instrument_id()
4769 ))
4770 .unwrap();
4771
4772 rust_strategy.py_register_indicator_for_quote_ticks(
4773 py,
4774 quote.instrument_id,
4775 indicator.clone_ref(py),
4776 );
4777 rust_strategy.py_register_indicator_for_trade_ticks(
4778 py,
4779 trade.instrument_id,
4780 indicator.clone_ref(py),
4781 );
4782 rust_strategy.py_register_indicator_for_bars(
4783 py,
4784 external_bar_type,
4785 indicator.clone_ref(py),
4786 );
4787
4788 DataActor::handle_quote(rust_strategy.inner_mut(), "e);
4789 DataActor::handle_trade(rust_strategy.inner_mut(), &trade);
4790 DataActor::handle_bar(rust_strategy.inner_mut(), &bar);
4791
4792 let events = events.extract::<Vec<String>>().unwrap();
4793
4794 assert_eq!(python_indicator_call_count(&indicator, py, "quote"), 1);
4795 assert_eq!(python_indicator_call_count(&indicator, py, "trade"), 1);
4796 assert_eq!(python_indicator_call_count(&indicator, py, "bar"), 1);
4797 assert_eq!(
4798 events,
4799 vec![
4800 "indicator:quote",
4801 "strategy:quote",
4802 "indicator:trade",
4803 "strategy:trade",
4804 "indicator:bar",
4805 "strategy:bar",
4806 ]
4807 );
4808 });
4809 }
4810
4811 #[rstest::rstest]
4812 fn test_registered_indicators_receive_historical_quote_trade_and_bar_batches() {
4813 pyo3::Python::initialize();
4814 Python::attach(|py| {
4815 let mut rust_strategy = PyStrategy::new(None);
4816 let indicator = create_tracking_python_indicator(py).unwrap();
4817 let quote = sample_quote();
4818 let trade = sample_trade();
4819 let bar = sample_bar();
4820 let quotes = vec![quote];
4821 let trades = vec![trade];
4822 let bars = vec![bar];
4823
4824 rust_strategy.py_register_indicator_for_quote_ticks(
4825 py,
4826 quote.instrument_id,
4827 indicator.clone_ref(py),
4828 );
4829 rust_strategy.py_register_indicator_for_trade_ticks(
4830 py,
4831 trade.instrument_id,
4832 indicator.clone_ref(py),
4833 );
4834 rust_strategy.py_register_indicator_for_bars(py, bar.bar_type, indicator.clone_ref(py));
4835
4836 let client_id = ClientId::new("TEST");
4837 let quotes_response = QuotesResponse::new(
4838 UUID4::new(),
4839 client_id,
4840 quote.instrument_id,
4841 quotes,
4842 None,
4843 None,
4844 UnixNanos::default(),
4845 None,
4846 );
4847 let trades_response = TradesResponse::new(
4848 UUID4::new(),
4849 client_id,
4850 trade.instrument_id,
4851 trades,
4852 None,
4853 None,
4854 UnixNanos::default(),
4855 None,
4856 );
4857 let bars_response = BarsResponse::new(
4858 UUID4::new(),
4859 client_id,
4860 bar.bar_type,
4861 bars,
4862 None,
4863 None,
4864 UnixNanos::default(),
4865 None,
4866 );
4867
4868 DataActor::handle_quotes_response(rust_strategy.inner_mut(), "es_response);
4869 DataActor::handle_trades_response(rust_strategy.inner_mut(), &trades_response);
4870 DataActor::handle_bars_response(rust_strategy.inner_mut(), &bars_response);
4871
4872 assert_eq!(python_indicator_call_count(&indicator, py, "quote"), 1);
4873 assert_eq!(python_indicator_call_count(&indicator, py, "trade"), 1);
4874 assert_eq!(python_indicator_call_count(&indicator, py, "bar"), 1);
4875 });
4876 }
4877
4878 fn assert_python_dispatch<F>(py: Python<'_>, method_name: &str, invoke: F) -> Py<PyAny>
4879 where
4880 F: FnOnce(&mut PyStrategy) -> anyhow::Result<()>,
4881 {
4882 let (py_strategy, mut rust_strategy) = create_registered_tracking_strategy(py);
4883 let result = invoke(&mut rust_strategy);
4884
4885 assert!(result.is_ok());
4886 assert!(python_method_was_called(&py_strategy, py, method_name));
4887 assert_eq!(python_method_call_count(&py_strategy, py, method_name), 1);
4888
4889 py_strategy
4890 }
4891
4892 fn assert_python_callback<F>(
4893 py: Python<'_>,
4894 method_name: &str,
4895 raises: bool,
4896 shutdown: bool,
4897 invoke: F,
4898 ) -> Py<PyAny>
4899 where
4900 F: Fn(&mut PyStrategy) -> anyhow::Result<()>,
4901 {
4902 let _guard = init_logging(
4903 TraderId::from("TRADER-001"),
4904 UUID4::new(),
4905 Default::default(),
4906 Default::default(),
4907 )
4908 .unwrap();
4909 let (py_strategy, mut rust_strategy) = create_registered_tracking_strategy(py);
4910 if raises {
4911 py_strategy.setattr(py, "raise_on", method_name).unwrap();
4912 }
4913
4914 arm_shutdown_on_error(shutdown);
4915 let result = invoke(&mut rust_strategy);
4916 let trigger = take_shutdown_on_error_trigger();
4917 disarm_shutdown_on_error();
4918
4919 assert!(result.is_ok());
4920 assert_eq!(python_method_call_count(&py_strategy, py, method_name), 1);
4921
4922 if raises && shutdown {
4923 let trigger = trigger.expect("Python callback failure must request shutdown");
4924 assert_eq!(
4925 trigger.component.as_str(),
4926 rust_strategy.inner().core.actor.actor_id.as_str()
4927 );
4928 assert!(
4929 trigger
4930 .message
4931 .contains(&format!("Python {method_name} failed:"))
4932 );
4933 assert!(trigger.message.contains("in fail_callback"));
4934 assert!(trigger.message.contains("ValueError: callback cause"));
4935 assert!(trigger.message.contains("RuntimeError: callback failure"));
4936 } else {
4937 assert_eq!(trigger, None);
4938 }
4939
4940 py_strategy
4941 }
4942
4943 #[rstest::rstest]
4944 #[case(false)]
4945 #[case(true)]
4946 fn test_python_timer_callback_exception_requests_shutdown(#[case] shutdown: bool) {
4947 Python::initialize();
4948 let _guard = init_logging(
4949 TraderId::from("TRADER-001"),
4950 UUID4::new(),
4951 Default::default(),
4952 Default::default(),
4953 )
4954 .unwrap();
4955 Python::attach(|py| {
4956 let (tracker, strategy) = create_registered_tracking_strategy(py);
4957 strategy.register_in_global_registries().unwrap();
4958 Component::start(strategy.inner_mut()).unwrap();
4959 tracker.setattr(py, "raise_on", "on_time_event").unwrap();
4960 let clock = DataActorNative::clock_rc(strategy.inner());
4961 clock
4962 .borrow_mut()
4963 .set_time_alert_ns("CALLBACK", UnixNanos::from(10), None, None)
4964 .unwrap();
4965 arm_shutdown_on_error(shutdown);
4966 let dispatched = dispatch_time_events(&clock, UnixNanos::from(10));
4967 let trigger = take_shutdown_on_error_trigger();
4968 disarm_shutdown_on_error();
4969
4970 assert_eq!(dispatched, 1);
4971 assert_eq!(python_method_call_count(&tracker, py, "on_time_event"), 1);
4972 assert_eq!(Component::state(strategy.inner()), ComponentState::Running);
4973
4974 if shutdown {
4975 let trigger = trigger.expect("strategy timer exception must request shutdown");
4976 assert!(trigger.message.contains("Python on_time_event failed:"));
4977 assert!(trigger.message.contains("in fail_callback"));
4978 assert!(trigger.message.contains("RuntimeError: callback failure"));
4979 } else {
4980 assert_eq!(trigger, None);
4981 }
4982 });
4983 }
4984
4985 #[rstest::rstest]
4986 fn test_python_data_callback_exception_preserves_traceback() {
4987 Python::initialize();
4988 Python::attach(|py| {
4989 let (tracker, strategy) = create_registered_tracking_strategy(py);
4990 tracker.setattr(py, "raise_on", "on_bar").unwrap();
4991 let error = DataActor::on_bar(strategy.inner_mut(), &sample_bar()).unwrap_err();
4992 let message = error.to_string();
4993
4994 assert_eq!(python_method_call_count(&tracker, py, "on_bar"), 1);
4995 assert!(message.contains("Python on_bar failed:"));
4996 assert!(message.contains("in fail_callback"));
4997 assert!(message.contains("ValueError: callback cause"));
4998 assert!(message.contains("RuntimeError: callback failure"));
4999 });
5000 }
5001
5002 #[rstest::rstest]
5003 #[case("on_market_exit")]
5004 #[case("post_market_exit")]
5005 fn test_python_dispatch_market_exit_callback_errors(
5006 #[case] method_name: &str,
5007 #[values(false, true)] raises: bool,
5008 #[values(false, true)] shutdown: bool,
5009 ) {
5010 Python::initialize();
5011 Python::attach(|py| {
5012 assert_python_callback(py, method_name, raises, shutdown, |strategy| {
5013 match method_name {
5014 "on_market_exit" => Strategy::on_market_exit(strategy.inner_mut()),
5015 "post_market_exit" => Strategy::post_market_exit(strategy.inner_mut()),
5016 _ => unreachable!(),
5017 }
5018
5019 Ok(())
5020 });
5021 });
5022 }
5023
5024 #[rstest::rstest]
5025 #[case("on_start")]
5026 #[case("on_stop")]
5027 #[case("on_resume")]
5028 #[case("on_reset")]
5029 #[case("on_dispose")]
5030 #[case("on_degrade")]
5031 #[case("on_fault")]
5032 fn test_python_dispatch_lifecycle_matrix(#[case] method_name: &str) {
5033 pyo3::Python::initialize();
5034 Python::attach(|py| {
5035 assert_python_dispatch(py, method_name, |rust_strategy| match method_name {
5036 "on_start" => DataActor::on_start(rust_strategy.inner_mut()),
5037 "on_stop" => DataActor::on_stop(rust_strategy.inner_mut()),
5038 "on_resume" => DataActor::on_resume(rust_strategy.inner_mut()),
5039 "on_reset" => DataActor::on_reset(rust_strategy.inner_mut()),
5040 "on_dispose" => DataActor::on_dispose(rust_strategy.inner_mut()),
5041 "on_degrade" => DataActor::on_degrade(rust_strategy.inner_mut()),
5042 "on_fault" => DataActor::on_fault(rust_strategy.inner_mut()),
5043 _ => unreachable!("unhandled lifecycle case: {method_name}"),
5044 });
5045 });
5046 }
5047
5048 #[rstest::rstest]
5049 #[case("on_save")]
5050 #[case("on_load")]
5051 fn test_python_dispatch_persistence_matrix(#[case] method_name: &str) {
5052 pyo3::Python::initialize();
5053 Python::attach(|py| {
5054 assert_python_dispatch(py, method_name, |rust_strategy| match method_name {
5055 "on_save" => {
5056 let state = DataActor::on_save(rust_strategy.inner()).unwrap();
5057 assert_eq!(
5058 state.get("strategy").map(Vec::as_slice),
5059 Some(b"saved".as_slice())
5060 );
5061 Ok(())
5062 }
5063 "on_load" => {
5064 let mut state = IndexMap::new();
5065 state.insert("strategy".to_string(), b"loaded".to_vec());
5066 DataActor::on_load(rust_strategy.inner_mut(), state)
5067 }
5068 _ => unreachable!("unhandled persistence case: {method_name}"),
5069 });
5070 });
5071 }
5072
5073 #[rstest::rstest]
5074 fn test_python_persistence_methods_convert_state() {
5075 pyo3::Python::initialize();
5076 Python::attach(|py| {
5077 let (py_strategy, mut rust_strategy) = create_registered_tracking_strategy(py);
5078
5079 let saved = rust_strategy.py_save(py).unwrap();
5080 let saved_state = saved
5081 .bind(py)
5082 .extract::<HashMap<String, Vec<u8>>>()
5083 .unwrap();
5084 assert_eq!(
5085 saved_state.get("strategy").map(Vec::as_slice),
5086 Some(&b"saved"[..])
5087 );
5088
5089 let load_state = PyDict::new(py);
5090 load_state
5091 .set_item("strategy", PyBytes::new(py, b"loaded-from-python"))
5092 .unwrap();
5093
5094 rust_strategy.py_load(&load_state).unwrap();
5095
5096 let loaded_state = python_last_loaded_state(&py_strategy, py).unwrap();
5097 assert_eq!(
5098 loaded_state.get("strategy").map(Vec::as_slice),
5099 Some(&b"loaded-from-python"[..])
5100 );
5101 });
5102 }
5103
5104 #[rstest::rstest]
5105 fn test_python_publish_data_and_signal_reach_msgbus() {
5106 use nautilus_common::msgbus::{
5107 MStr, MessageBus, Pattern, get_message_bus, switchboard::get_custom_topic,
5108 typed_handler::ShareableMessageHandler,
5109 };
5110 use nautilus_core::python::IntoPyObjectNautilusExt;
5111
5112 *get_message_bus().borrow_mut() = MessageBus::default();
5113
5114 pyo3::Python::initialize();
5115 Python::attach(|py| {
5116 let (_, rust_strategy) = create_registered_tracking_strategy(py);
5117 let data = sample_data();
5118
5119 let received_data: Rc<RefCell<Vec<CustomData>>> = Rc::new(RefCell::new(Vec::new()));
5120 let received_data_clone = received_data.clone();
5121 let data_handler = ShareableMessageHandler::from_typed(move |data: &CustomData| {
5122 received_data_clone.borrow_mut().push(data.clone());
5123 });
5124 msgbus::subscribe_any(get_custom_topic(&data.data_type).into(), data_handler, None);
5125
5126 let received_signals: Rc<RefCell<Vec<Signal>>> = Rc::new(RefCell::new(Vec::new()));
5127 let received_signals_clone = received_signals.clone();
5128 let signal_handler = ShareableMessageHandler::from_typed(move |data: &CustomData| {
5129 if let Some(signal) = data.data.as_any().downcast_ref::<Signal>() {
5130 received_signals_clone.borrow_mut().push(signal.clone());
5131 }
5132 });
5133 let signal_pattern: MStr<Pattern> = "data.Signal*".to_string().into();
5134 msgbus::subscribe_any(signal_pattern, signal_handler, None);
5135
5136 rust_strategy
5137 .py_publish_data(&data.data_type, &data)
5138 .unwrap();
5139
5140 let value: Py<PyAny> = 2.0_f64.into_py_any_unwrap(py);
5141 rust_strategy
5142 .py_publish_signal(py, "risk", value, 1_700_000_000_000_000_000)
5143 .unwrap();
5144
5145 let received_data = received_data.borrow();
5146 assert_eq!(received_data.len(), 1);
5147 assert_eq!(received_data[0].data_type, data.data_type);
5148
5149 let received_signals = received_signals.borrow();
5150 assert_eq!(received_signals.len(), 1);
5151 assert_eq!(received_signals[0].name, "risk");
5152 assert_eq!(received_signals[0].value, "2.0");
5153 assert_eq!(
5154 received_signals[0].ts_event,
5155 UnixNanos::from(1_700_000_000_000_000_000_u64),
5156 );
5157 });
5158 }
5159
5160 #[rstest::rstest]
5161 fn test_python_add_and_update_synthetic_update_cache() {
5162 use std::str::FromStr;
5163
5164 use nautilus_model::{
5165 identifiers::{InstrumentId, Symbol},
5166 instruments::SyntheticInstrument,
5167 };
5168
5169 pyo3::Python::initialize();
5170 Python::attach(|py| {
5171 let (_, rust_strategy) = create_registered_tracking_strategy(py);
5172
5173 let comp1 = InstrumentId::from_str("BTC-USD.VENUE").unwrap();
5174 let comp2 = InstrumentId::from_str("ETH-USD.VENUE").unwrap();
5175 let symbol = Symbol::from("SYN");
5176 let original_formula = format!("({comp1} + {comp2}) / 2.0");
5177 let synthetic = SyntheticInstrument::builder()
5178 .symbol(symbol)
5179 .price_precision(2)
5180 .components(vec![comp1, comp2])
5181 .formula(&original_formula)
5182 .ts_event(UnixNanos::default())
5183 .ts_init(UnixNanos::default())
5184 .build()
5185 .unwrap();
5186 let synthetic_id = synthetic.id;
5187
5188 rust_strategy.py_add_synthetic(synthetic).unwrap();
5189
5190 let updated_formula = format!("{comp1} + {comp2}");
5191 let updated = SyntheticInstrument::builder()
5192 .symbol(symbol)
5193 .price_precision(2)
5194 .components(vec![comp1, comp2])
5195 .formula(&updated_formula)
5196 .ts_event(UnixNanos::default())
5197 .ts_init(UnixNanos::default())
5198 .build()
5199 .unwrap();
5200 rust_strategy.py_update_synthetic(updated).unwrap();
5201
5202 let cache = DataActor::cache(rust_strategy.inner());
5203 let stored = cache.synthetic(&synthetic_id).unwrap();
5204 assert_eq!(stored.formula, updated_formula);
5205 });
5206 }
5207
5208 #[rstest::rstest]
5209 fn test_python_subscribe_and_unsubscribe_signal_update_msgbus() {
5210 use nautilus_common::msgbus::{MessageBus, get_message_bus, switchboard::get_signal_topic};
5211
5212 *get_message_bus().borrow_mut() = MessageBus::default();
5213
5214 pyo3::Python::initialize();
5215 Python::attach(|py| {
5216 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
5217
5218 rust_strategy.py_subscribe_signal("risk", Some(50)).unwrap();
5219
5220 let topic = get_signal_topic("risk");
5221 let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
5222 assert_eq!(subscriptions.len(), 1);
5223 assert_eq!(subscriptions[0].priority, 50);
5224
5225 rust_strategy.py_unsubscribe_signal("risk").unwrap();
5226
5227 let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
5228 assert!(subscriptions.is_empty());
5229 });
5230 }
5231
5232 #[rstest::rstest]
5233 #[case(None)]
5234 #[case(Some(SystemChannel::ExecCommands))]
5235 fn test_python_subscribe_and_unsubscribe_queue_state_update_msgbus(
5236 #[case] channel: Option<SystemChannel>,
5237 ) {
5238 use nautilus_common::msgbus::{MessageBus, MessagingSwitchboard, get_message_bus};
5239
5240 *get_message_bus().borrow_mut() = MessageBus::default();
5241
5242 pyo3::Python::initialize();
5243 Python::attach(|py| {
5244 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
5245
5246 rust_strategy
5247 .py_subscribe_queue_state(channel, Some(50))
5248 .unwrap();
5249
5250 let topic =
5251 MessagingSwitchboard::queue_state_changed_topic(SystemChannel::ExecCommands);
5252 let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
5253 assert_eq!(subscriptions.len(), 1);
5254 assert_eq!(subscriptions[0].priority, 50);
5255 let unrelated =
5256 MessagingSwitchboard::queue_state_changed_topic(SystemChannel::DataEvents);
5257 assert_eq!(
5258 get_message_bus()
5259 .borrow_mut()
5260 .matching_subscriptions(unrelated)
5261 .len(),
5262 usize::from(channel.is_none())
5263 );
5264
5265 rust_strategy.py_unsubscribe_queue_state(channel).unwrap();
5266
5267 let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
5268 assert!(subscriptions.is_empty());
5269 });
5270 }
5271
5272 #[rstest::rstest]
5273 #[case(None, None)]
5274 #[case(Some(ClientId::from("BINANCE")), None)]
5275 #[case(None, Some("binance-futures-market-streams"))]
5276 #[case(
5277 Some(ClientId::from("BINANCE")),
5278 Some("binance-futures-market-streams")
5279 )]
5280 fn test_python_subscribe_and_unsubscribe_socket_state_update_msgbus(
5281 #[case] client_id: Option<ClientId>,
5282 #[case] endpoint: Option<&str>,
5283 ) {
5284 use nautilus_common::msgbus::{MessageBus, MessagingSwitchboard, get_message_bus};
5285
5286 *get_message_bus().borrow_mut() = MessageBus::default();
5287
5288 pyo3::Python::initialize();
5289 Python::attach(|py| {
5290 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
5291
5292 rust_strategy
5293 .py_subscribe_socket_state(client_id, endpoint, Some(50))
5294 .unwrap();
5295
5296 let topic = MessagingSwitchboard::socket_state_changed_topic(
5297 ClientId::from("BINANCE"),
5298 "binance-futures-market-streams",
5299 );
5300 let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
5301 assert_eq!(subscriptions.len(), 1);
5302 assert_eq!(subscriptions[0].priority, 50);
5303 let unrelated =
5304 MessagingSwitchboard::socket_state_changed_topic(ClientId::from("BYBIT"), "orders");
5305 assert_eq!(
5306 get_message_bus()
5307 .borrow_mut()
5308 .matching_subscriptions(unrelated)
5309 .len(),
5310 usize::from(client_id.is_none() && endpoint.is_none())
5311 );
5312
5313 rust_strategy
5314 .py_unsubscribe_socket_state(client_id, endpoint)
5315 .unwrap();
5316
5317 let subscriptions = get_message_bus().borrow_mut().matching_subscriptions(topic);
5318 assert!(subscriptions.is_empty());
5319 });
5320 }
5321
5322 #[rstest::rstest]
5323 fn test_python_reconnect_socket_enqueues_typed_command() {
5324 pyo3::Python::initialize();
5325 Python::attach(|py| {
5326 let (system_tx, mut system_rx) = tokio::sync::mpsc::unbounded_channel();
5327 replace_system_command_sender(system_tx);
5328 let (_, strategy) = create_registered_tracking_strategy(py);
5329
5330 strategy
5331 .py_reconnect_socket(ClientId::from("POLYMARKET"), "polymarket-market-streams")
5332 .expect("valid reconnect command");
5333 let command = system_rx
5334 .try_recv()
5335 .expect("reconnect command should be queued");
5336 let SystemCommand::ReconnectSocket(command) = command;
5337
5338 assert_eq!(command.trader_id, TraderId::from("TRADER-001"));
5339 assert_eq!(command.client_id, ClientId::from("POLYMARKET"));
5340 assert_eq!(command.endpoint, "polymarket-market-streams");
5341 assert_eq!(command.ts_init, UnixNanos::default());
5342 });
5343 }
5344
5345 #[rstest::rstest]
5346 fn test_python_stop_stops_immediately_when_manage_stop_disabled() {
5347 pyo3::Python::initialize();
5348 Python::attach(|py| {
5349 let config = StrategyConfig {
5350 strategy_id: Some(StrategyId::from("TEST-001")),
5351 order_id_tag: Some("001".to_string()),
5352 manage_stop: false,
5353 ..Default::default()
5354 };
5355 let (py_strategy, mut rust_strategy) =
5356 create_registered_tracking_strategy_with_config(py, Some(config));
5357
5358 rust_strategy.py_start().unwrap();
5359 rust_strategy.py_stop().unwrap();
5360
5361 assert!(rust_strategy.py_is_stopped());
5362 assert!(!rust_strategy.inner().core.pending_stop);
5363 assert!(!rust_strategy.inner().core.is_exiting);
5364 assert_eq!(python_method_call_count(&py_strategy, py, "on_stop"), 1);
5365 });
5366 }
5367
5368 #[rstest::rstest]
5369 fn test_python_stop_defers_when_manage_stop_enabled() {
5370 pyo3::Python::initialize();
5371 Python::attach(|py| {
5372 let config = StrategyConfig {
5373 strategy_id: Some(StrategyId::from("TEST-001")),
5374 order_id_tag: Some("001".to_string()),
5375 manage_stop: true,
5376 ..Default::default()
5377 };
5378 let (py_strategy, mut rust_strategy) =
5379 create_registered_tracking_strategy_with_config(py, Some(config));
5380
5381 rust_strategy.py_start().unwrap();
5382 rust_strategy.py_stop().unwrap();
5383
5384 assert!(rust_strategy.py_is_running());
5385 assert!(rust_strategy.inner().core.pending_stop);
5386 assert!(rust_strategy.inner().core.is_exiting);
5387 assert_eq!(python_method_call_count(&py_strategy, py, "on_stop"), 0);
5388 });
5389 }
5390
5391 #[rstest::rstest]
5392 fn test_python_market_exit_methods_update_state_and_dispatch_hooks() {
5393 pyo3::Python::initialize();
5394 Python::attach(|py| {
5395 let (py_strategy, mut rust_strategy) = create_registered_tracking_strategy(py);
5396
5397 rust_strategy.py_start().unwrap();
5398
5399 assert!(!rust_strategy.py_is_exiting());
5400
5401 rust_strategy.py_market_exit().unwrap();
5402
5403 assert!(rust_strategy.py_is_exiting());
5404 assert_eq!(
5405 python_method_call_count(&py_strategy, py, "on_market_exit"),
5406 1
5407 );
5408
5409 rust_strategy.inner_mut().finalize_market_exit();
5410
5411 assert!(!rust_strategy.py_is_exiting());
5412 assert_eq!(
5413 python_method_call_count(&py_strategy, py, "post_market_exit"),
5414 1
5415 );
5416 });
5417 }
5418
5419 #[rstest::rstest]
5420 #[case::order_list_object(true)]
5421 #[case::raw_order_sequence(false)]
5422 fn test_python_submit_order_list_accepts_order_list_inputs(#[case] wrap_order_list: bool) {
5423 pyo3::Python::initialize();
5424 Python::attach(|py| {
5425 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
5426 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5427 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
5428 msgbus::register_trading_command_endpoint(
5429 MessagingSwitchboard::risk_engine_queue_execute(),
5430 risk_handler,
5431 );
5432
5433 let strategy_id = rust_strategy.strategy_id();
5434 let client_order_id1 = ClientOrderId::from("O-PYO3-LIST-001");
5435 let client_order_id2 = ClientOrderId::from("O-PYO3-LIST-002");
5436 let orders = vec![
5437 sample_python_market_order(py, strategy_id, client_order_id1).unwrap(),
5438 sample_python_market_order(py, strategy_id, client_order_id2).unwrap(),
5439 ];
5440 let params = PyDict::new(py);
5441
5442 params.set_item("routing_hint", "prefer_batch").unwrap();
5443 let order_list = if wrap_order_list {
5444 let order_list_type = py
5445 .eval(c_str!("type('OrderListShim', (), {})"), None, None)
5446 .unwrap();
5447 let order_list = order_list_type.call0().unwrap();
5448
5449 order_list.setattr("orders", orders).unwrap();
5450 order_list.unbind()
5451 } else {
5452 PyList::new(py, orders).unwrap().into_any().unbind()
5453 };
5454
5455 rust_strategy
5456 .py_submit_order_list(py, order_list, None, None, Some(params.unbind()))
5457 .unwrap();
5458
5459 let cache = DataActor::cache(rust_strategy.inner());
5460 let cached_order1 = cache.order(&client_order_id1).unwrap();
5461 let cached_order2 = cache.order(&client_order_id2).unwrap();
5462 let order_list_id = cached_order1.order_list_id().unwrap();
5463 let order_list = cache.order_list(&order_list_id).unwrap();
5464
5465 assert_eq!(cached_order2.order_list_id(), Some(order_list_id));
5466 assert_eq!(
5467 order_list.client_order_ids.as_slice(),
5468 &[client_order_id1, client_order_id2]
5469 );
5470
5471 let py_cache =
5472 Py::new(py, PyCache::from_rc(rust_strategy.inner().core.cache_rc())).unwrap();
5473 let py_order_list = py_cache
5474 .bind(py)
5475 .call_method1("order_list", (order_list_id,))
5476 .unwrap();
5477
5478 assert_eq!(
5479 py_order_list
5480 .getattr("id")
5481 .unwrap()
5482 .extract::<OrderListId>()
5483 .unwrap(),
5484 order_list_id,
5485 );
5486 assert_eq!(
5487 py_order_list
5488 .call_method0("client_order_ids")
5489 .unwrap()
5490 .extract::<Vec<ClientOrderId>>()
5491 .unwrap(),
5492 vec![client_order_id1, client_order_id2],
5493 );
5494
5495 let risk_messages = risk_messages.get_messages();
5496 assert_eq!(risk_messages.len(), 1);
5497 let Some(TradingCommand::SubmitOrderList(command)) = risk_messages.first() else {
5498 panic!("expected SubmitOrderList command");
5499 };
5500 assert_eq!(
5501 command
5502 .params
5503 .as_ref()
5504 .and_then(|params| params.get("routing_hint")),
5505 Some(&Value::String("prefer_batch".to_string()))
5506 );
5507 });
5508 }
5509
5510 #[rstest::rstest]
5511 fn test_python_close_position_forwards_params_to_submit_order() {
5512 pyo3::Python::initialize();
5513 Python::attach(|py| {
5514 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
5515 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5516 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
5517 msgbus::register_trading_command_endpoint(
5518 MessagingSwitchboard::risk_engine_queue_execute(),
5519 risk_handler,
5520 );
5521
5522 let position_id = PositionId::from("P-PYO3-CLOSE-001");
5523 let position = sample_open_position(
5524 rust_strategy.strategy_id(),
5525 position_id,
5526 ClientOrderId::from("O-PYO3-CLOSE-001"),
5527 );
5528 let params = PyDict::new(py);
5529
5530 params.set_item("routing_hint", "close_single").unwrap();
5531 rust_strategy
5532 .py_close_position(
5533 &position,
5534 None,
5535 None,
5536 None,
5537 None,
5538 None,
5539 Some(params.unbind()),
5540 )
5541 .unwrap();
5542
5543 let risk_messages = risk_messages.get_messages();
5544 assert_eq!(risk_messages.len(), 1);
5545 let Some(TradingCommand::SubmitOrder(command)) = risk_messages.first() else {
5546 panic!("expected SubmitOrder command");
5547 };
5548 assert_eq!(command.position_id, Some(position_id));
5549 assert_eq!(
5550 command
5551 .params
5552 .as_ref()
5553 .and_then(|params| params.get("routing_hint")),
5554 Some(&Value::String("close_single".to_string()))
5555 );
5556 });
5557 }
5558
5559 #[rstest::rstest]
5560 fn test_python_close_all_positions_forwards_params_to_submit_order() {
5561 pyo3::Python::initialize();
5562 Python::attach(|py| {
5563 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
5564 let (risk_handler, risk_messages): (_, TypedIntoMessageSavingHandler<TradingCommand>) =
5565 get_typed_into_message_saving_handler(Some(Ustr::from("RiskEngine.queue_execute")));
5566 msgbus::register_trading_command_endpoint(
5567 MessagingSwitchboard::risk_engine_queue_execute(),
5568 risk_handler,
5569 );
5570
5571 let instrument = sample_instrument();
5572 let position_id1 = PositionId::from("P-PYO3-CLOSE-ALL-001");
5573 let position_id2 = PositionId::from("P-PYO3-CLOSE-ALL-002");
5574 let position1 = sample_open_position(
5575 rust_strategy.strategy_id(),
5576 position_id1,
5577 ClientOrderId::from("O-PYO3-CLOSE-ALL-001"),
5578 );
5579 let position2 = sample_open_position(
5580 rust_strategy.strategy_id(),
5581 position_id2,
5582 ClientOrderId::from("O-PYO3-CLOSE-ALL-002"),
5583 );
5584 let cache = rust_strategy.inner().core.cache_rc();
5585
5586 {
5587 let mut cache = cache.borrow_mut();
5588 cache
5589 .add_instrument(InstrumentAny::CurrencyPair(instrument.clone()))
5590 .unwrap();
5591 cache.add_position(&position1, OmsType::Hedging).unwrap();
5592 cache.add_position(&position2, OmsType::Hedging).unwrap();
5593 }
5594
5595 let params = PyDict::new(py);
5596
5597 params.set_item("routing_hint", "close_all").unwrap();
5598 rust_strategy
5599 .py_close_all_positions(
5600 instrument.id,
5601 None,
5602 None,
5603 None,
5604 None,
5605 None,
5606 None,
5607 Some(params.unbind()),
5608 )
5609 .unwrap();
5610
5611 let risk_messages = risk_messages.get_messages();
5612 assert_eq!(risk_messages.len(), 2);
5613 let commands: Vec<_> = risk_messages
5614 .iter()
5615 .map(|message| {
5616 let TradingCommand::SubmitOrder(command) = message else {
5617 panic!("expected SubmitOrder command");
5618 };
5619 command
5620 })
5621 .collect();
5622
5623 assert!(
5624 commands
5625 .iter()
5626 .any(|command| command.position_id == Some(position_id1))
5627 );
5628 assert!(
5629 commands
5630 .iter()
5631 .any(|command| command.position_id == Some(position_id2))
5632 );
5633 assert!(commands.iter().all(|command| {
5634 command
5635 .params
5636 .as_ref()
5637 .and_then(|params| params.get("routing_hint"))
5638 == Some(&Value::String("close_all".to_string()))
5639 }));
5640 });
5641 }
5642
5643 #[rstest::rstest]
5644 fn test_python_cancel_gtd_expiry_accepts_order() {
5645 pyo3::Python::initialize();
5646 Python::attach(|py| {
5647 let (_, mut rust_strategy) = create_registered_tracking_strategy(py);
5648 let strategy_id = rust_strategy.strategy_id();
5649 let client_order_id = ClientOrderId::from("O-PYO3-GTD-001");
5650 let timer_name = format!("GTD-EXPIRY:{client_order_id}");
5651 let order = OrderTestBuilder::new(OrderType::Limit)
5652 .trader_id(TraderId::from("TRADER-001"))
5653 .strategy_id(strategy_id)
5654 .instrument_id(sample_instrument().id)
5655 .client_order_id(client_order_id)
5656 .quantity(Quantity::from(100_000))
5657 .price(Price::from("1.00000"))
5658 .time_in_force(TimeInForce::Gtd)
5659 .expire_time(UnixNanos::from(1))
5660 .build();
5661 let py_order = order_any_to_pyobject(py, order).unwrap();
5662
5663 {
5664 let mut clock = rust_strategy.inner_mut().core.clock_mut();
5665 clock
5666 .set_time_alert_ns(&timer_name, UnixNanos::from(1), None, None)
5667 .unwrap();
5668 }
5669 rust_strategy
5670 .inner_mut()
5671 .core
5672 .gtd_timers
5673 .insert(client_order_id, Ustr::from(&timer_name));
5674
5675 rust_strategy
5676 .py_cancel_gtd_expiry(py, py_order)
5677 .expect("cancel_gtd_expiry should accept Python order");
5678
5679 let clock_timer_exists = rust_strategy
5680 .inner_mut()
5681 .core
5682 .clock_mut()
5683 .timer_names()
5684 .contains(&timer_name.as_str());
5685
5686 assert!(
5687 !rust_strategy
5688 .inner_mut()
5689 .has_gtd_expiry_timer(&client_order_id)
5690 );
5691 assert!(!clock_timer_exists);
5692 });
5693 }
5694
5695 #[rstest::rstest]
5696 #[case("on_time_event")]
5697 #[case("on_data")]
5698 #[case("on_signal")]
5699 #[case("on_queue_state")]
5700 #[case("on_socket_state")]
5701 #[case("on_instrument")]
5702 #[case("on_quote")]
5703 #[case("on_trade")]
5704 #[case("on_bar")]
5705 #[case("on_book_deltas")]
5706 #[case("on_book_depth")]
5707 #[case("on_book")]
5708 #[case("on_mark_price")]
5709 #[case("on_index_price")]
5710 #[case("on_funding_rate")]
5711 #[case("on_instrument_status")]
5712 #[case("on_instrument_close")]
5713 #[case("on_option_greeks")]
5714 #[case("on_option_chain")]
5715 #[case("on_historical_data")]
5716 #[case("on_historical_book_deltas")]
5717 #[case("on_historical_book_depth")]
5718 #[case("on_historical_quotes")]
5719 #[case("on_historical_trades")]
5720 #[case("on_historical_funding_rates")]
5721 #[case("on_historical_bars")]
5722 #[case("on_historical_mark_prices")]
5723 #[case("on_historical_index_prices")]
5724 fn test_python_dispatch_data_callback_matrix(#[case] method_name: &str) {
5725 pyo3::Python::initialize();
5726 Python::attach(|py| {
5727 assert_python_dispatch(py, method_name, |rust_strategy| match method_name {
5728 "on_time_event" => {
5729 let event = sample_time_event();
5730 DataActor::on_time_event(rust_strategy.inner_mut(), &event)
5731 }
5732 "on_data" => {
5733 let data = sample_data();
5734 rust_strategy.inner_mut().on_data(&data)
5735 }
5736 "on_signal" => {
5737 let signal = sample_signal();
5738 rust_strategy.inner_mut().on_signal(&signal)
5739 }
5740 "on_queue_state" => {
5741 let event = sample_queue_state_changed();
5742 rust_strategy.inner_mut().on_queue_state(&event)
5743 }
5744 "on_socket_state" => {
5745 let event = sample_socket_state_changed();
5746 rust_strategy.inner_mut().on_socket_state(&event)
5747 }
5748 "on_instrument" => {
5749 let instrument = InstrumentAny::CurrencyPair(sample_instrument());
5750 rust_strategy.inner_mut().on_instrument(&instrument)
5751 }
5752 "on_quote" => {
5753 let quote = sample_quote();
5754 rust_strategy.inner_mut().on_quote("e)
5755 }
5756 "on_trade" => {
5757 let trade = sample_trade();
5758 rust_strategy.inner_mut().on_trade(&trade)
5759 }
5760 "on_bar" => {
5761 let bar = sample_bar();
5762 rust_strategy.inner_mut().on_bar(&bar)
5763 }
5764 "on_book_deltas" => {
5765 let deltas = sample_book_deltas();
5766 rust_strategy.inner_mut().on_book_deltas(&deltas)
5767 }
5768 "on_book_depth" => {
5769 let depth = sample_book_depth();
5770 rust_strategy.inner_mut().on_book_depth(&depth)
5771 }
5772 "on_book" => {
5773 let book = sample_book();
5774 rust_strategy.inner_mut().on_book(&book)
5775 }
5776 "on_mark_price" => {
5777 let mark_price = sample_mark_price();
5778 rust_strategy.inner_mut().on_mark_price(&mark_price)
5779 }
5780 "on_index_price" => {
5781 let index_price = sample_index_price();
5782 rust_strategy.inner_mut().on_index_price(&index_price)
5783 }
5784 "on_funding_rate" => {
5785 let funding_rate = sample_funding_rate();
5786 rust_strategy.inner_mut().on_funding_rate(&funding_rate)
5787 }
5788 "on_instrument_status" => {
5789 let status = sample_instrument_status();
5790 rust_strategy.inner_mut().on_instrument_status(&status)
5791 }
5792 "on_instrument_close" => {
5793 let close = sample_instrument_close();
5794 rust_strategy.inner_mut().on_instrument_close(&close)
5795 }
5796 "on_option_greeks" => {
5797 let greeks = sample_option_greeks();
5798 DataActor::on_option_greeks(rust_strategy.inner_mut(), &greeks)
5799 }
5800 "on_option_chain" => {
5801 let slice = sample_option_chain();
5802 DataActor::on_option_chain(rust_strategy.inner_mut(), &slice)
5803 }
5804 "on_historical_data" => {
5805 let data = sample_data();
5806 rust_strategy.inner_mut().on_historical_data(&data)
5807 }
5808 "on_historical_book_deltas" => {
5809 let deltas = sample_book_deltas().deltas;
5810 rust_strategy.inner_mut().on_historical_book_deltas(&deltas)
5811 }
5812 "on_historical_book_depth" => {
5813 let depths = vec![sample_book_depth()];
5814 rust_strategy.inner_mut().on_historical_book_depth(&depths)
5815 }
5816 "on_historical_quotes" => {
5817 let quotes = vec![sample_quote()];
5818 rust_strategy.inner_mut().on_historical_quotes("es)
5819 }
5820 "on_historical_trades" => {
5821 let trades = vec![sample_trade()];
5822 rust_strategy.inner_mut().on_historical_trades(&trades)
5823 }
5824 "on_historical_funding_rates" => {
5825 let funding_rates = vec![sample_funding_rate()];
5826 rust_strategy
5827 .inner_mut()
5828 .on_historical_funding_rates(&funding_rates)
5829 }
5830 "on_historical_bars" => {
5831 let bars = vec![sample_bar()];
5832 rust_strategy.inner_mut().on_historical_bars(&bars)
5833 }
5834 "on_historical_mark_prices" => {
5835 let mark_prices = vec![sample_mark_price()];
5836 rust_strategy
5837 .inner_mut()
5838 .on_historical_mark_prices(&mark_prices)
5839 }
5840 "on_historical_index_prices" => {
5841 let index_prices = vec![sample_index_price()];
5842 rust_strategy
5843 .inner_mut()
5844 .on_historical_index_prices(&index_prices)
5845 }
5846 _ => unreachable!("unhandled data callback case: {method_name}"),
5847 });
5848 });
5849 }
5850
5851 #[rstest::rstest]
5852 fn test_python_dispatch_historical_custom_data_preserves_payload_shape() {
5853 pyo3::Python::initialize();
5854
5855 Python::attach(|py| {
5856 let scalar = stub_custom_data(3, 126, None, None);
5857 let py_strategy = assert_python_dispatch(py, "on_historical_data", |rust_strategy| {
5858 rust_strategy.inner_mut().on_historical_data(&scalar)
5859 });
5860 let actual_scalar = py_strategy
5861 .call_method1(py, "last_call_args", ("on_historical_data",))
5862 .unwrap()
5863 .bind(py)
5864 .get_item(0)
5865 .unwrap()
5866 .extract::<CustomData>()
5867 .unwrap();
5868
5869 assert_eq!(actual_scalar, scalar);
5870
5871 let expected = vec![
5872 stub_custom_data(1, 42, None, None),
5873 stub_custom_data(2, 84, None, None),
5874 ];
5875 let py_strategy = assert_python_dispatch(py, "on_historical_data", |rust_strategy| {
5876 rust_strategy.inner_mut().on_historical_data(&expected)
5877 });
5878 let actual = py_strategy
5879 .call_method1(py, "last_call_args", ("on_historical_data",))
5880 .unwrap()
5881 .bind(py)
5882 .get_item(0)
5883 .unwrap()
5884 .extract::<Vec<CustomData>>()
5885 .unwrap();
5886
5887 assert_eq!(actual, expected);
5888 });
5889 }
5890
5891 #[rstest::rstest]
5892 fn test_python_dispatch_historical_book_deltas_preserves_batch() {
5893 pyo3::Python::initialize();
5894
5895 Python::attach(|py| {
5896 let expected = stub_deltas().deltas;
5897 let py_strategy =
5898 assert_python_dispatch(py, "on_historical_book_deltas", |rust_strategy| {
5899 rust_strategy
5900 .inner_mut()
5901 .on_historical_book_deltas(&expected)
5902 });
5903 let actual = py_strategy
5904 .call_method1(py, "last_call_args", ("on_historical_book_deltas",))
5905 .unwrap()
5906 .bind(py)
5907 .get_item(0)
5908 .unwrap()
5909 .extract::<Vec<OrderBookDelta>>()
5910 .unwrap();
5911
5912 assert_eq!(actual, expected);
5913 });
5914 }
5915
5916 #[rstest::rstest]
5917 fn test_python_dispatch_historical_book_depth_preserves_batch() {
5918 pyo3::Python::initialize();
5919
5920 Python::attach(|py| {
5921 let first = stub_depth10();
5922 let mut second = first.clone();
5923 second.sequence = 17;
5924 second.ts_event = UnixNanos::from(18);
5925 second.ts_init = UnixNanos::from(19);
5926 let expected = vec![first, second];
5927 let py_strategy =
5928 assert_python_dispatch(py, "on_historical_book_depth", |rust_strategy| {
5929 rust_strategy
5930 .inner_mut()
5931 .on_historical_book_depth(&expected)
5932 });
5933 let actual = py_strategy
5934 .call_method1(py, "last_call_args", ("on_historical_book_depth",))
5935 .unwrap()
5936 .bind(py)
5937 .get_item(0)
5938 .unwrap()
5939 .extract::<Vec<OrderBookDepth>>()
5940 .unwrap();
5941
5942 assert_eq!(actual, expected);
5943 });
5944 }
5945
5946 #[rstest::rstest]
5947 #[case("on_order_initialized")]
5948 #[case("on_order_event")]
5949 #[case("on_order_denied")]
5950 #[case("on_order_emulated")]
5951 #[case("on_order_released")]
5952 #[case("on_order_submitted")]
5953 #[case("on_order_rejected")]
5954 #[case("on_order_accepted")]
5955 #[case("on_order_expired")]
5956 #[case("on_order_triggered")]
5957 #[case("on_order_pending_update")]
5958 #[case("on_order_pending_cancel")]
5959 #[case("on_order_modify_rejected")]
5960 #[case("on_order_cancel_rejected")]
5961 #[case("on_order_updated")]
5962 #[case("on_order_canceled")]
5963 #[case("on_order_filled")]
5964 #[case("on_order_fill_voided")]
5965 fn test_python_dispatch_order_callback_matrix(
5966 #[case] method_name: &str,
5967 #[values(false, true)] raises: bool,
5968 #[values(false, true)] shutdown: bool,
5969 ) {
5970 pyo3::Python::initialize();
5971 Python::attach(|py| {
5972 assert_python_callback(py, method_name, raises, shutdown, |rust_strategy| {
5973 match method_name {
5974 "on_order_initialized" => {
5975 Strategy::on_order_initialized(
5976 rust_strategy.inner_mut(),
5977 OrderInitialized::default(),
5978 );
5979 Ok(())
5980 }
5981 "on_order_event" => {
5982 Strategy::on_order_event(
5983 rust_strategy.inner_mut(),
5984 OrderEventAny::Accepted(OrderAccepted::default()),
5985 );
5986 Ok(())
5987 }
5988 "on_order_denied" => {
5989 Strategy::on_order_denied(
5990 rust_strategy.inner_mut(),
5991 OrderDenied::default(),
5992 );
5993 Ok(())
5994 }
5995 "on_order_emulated" => {
5996 Strategy::on_order_emulated(
5997 rust_strategy.inner_mut(),
5998 OrderEmulated::default(),
5999 );
6000 Ok(())
6001 }
6002 "on_order_released" => {
6003 Strategy::on_order_released(
6004 rust_strategy.inner_mut(),
6005 OrderReleased::default(),
6006 );
6007 Ok(())
6008 }
6009 "on_order_submitted" => {
6010 Strategy::on_order_submitted(
6011 rust_strategy.inner_mut(),
6012 OrderSubmitted::default(),
6013 );
6014 Ok(())
6015 }
6016 "on_order_rejected" => {
6017 Strategy::on_order_rejected(
6018 rust_strategy.inner_mut(),
6019 OrderRejected::default(),
6020 );
6021 Ok(())
6022 }
6023 "on_order_accepted" => {
6024 Strategy::on_order_accepted(
6025 rust_strategy.inner_mut(),
6026 OrderAccepted::default(),
6027 );
6028 Ok(())
6029 }
6030 "on_order_expired" => {
6031 Strategy::on_order_expired(
6032 rust_strategy.inner_mut(),
6033 OrderExpired::default(),
6034 );
6035 Ok(())
6036 }
6037 "on_order_triggered" => {
6038 Strategy::on_order_triggered(
6039 rust_strategy.inner_mut(),
6040 OrderTriggered::default(),
6041 );
6042 Ok(())
6043 }
6044 "on_order_pending_update" => {
6045 Strategy::on_order_pending_update(
6046 rust_strategy.inner_mut(),
6047 OrderPendingUpdate::default(),
6048 );
6049 Ok(())
6050 }
6051 "on_order_pending_cancel" => {
6052 Strategy::on_order_pending_cancel(
6053 rust_strategy.inner_mut(),
6054 OrderPendingCancel::default(),
6055 );
6056 Ok(())
6057 }
6058 "on_order_modify_rejected" => {
6059 Strategy::on_order_modify_rejected(
6060 rust_strategy.inner_mut(),
6061 OrderModifyRejected::default(),
6062 );
6063 Ok(())
6064 }
6065 "on_order_cancel_rejected" => {
6066 Strategy::on_order_cancel_rejected(
6067 rust_strategy.inner_mut(),
6068 OrderCancelRejected::default(),
6069 );
6070 Ok(())
6071 }
6072 "on_order_updated" => {
6073 Strategy::on_order_updated(
6074 rust_strategy.inner_mut(),
6075 OrderUpdated::default(),
6076 );
6077 Ok(())
6078 }
6079 "on_order_canceled" => {
6080 let event = OrderCanceled::default();
6081 Strategy::on_order_canceled(rust_strategy.inner_mut(), &event);
6082 Ok(())
6083 }
6084 "on_order_filled" => {
6085 let event = OrderFilledSpec::builder().build();
6086 Strategy::on_order_filled(rust_strategy.inner_mut(), &event);
6087 Ok(())
6088 }
6089 "on_order_fill_voided" => {
6090 let event = OrderFillVoidedSpec::builder().build();
6091 Strategy::on_order_fill_voided(rust_strategy.inner_mut(), &event);
6092 Ok(())
6093 }
6094 _ => unreachable!("unhandled order callback case: {method_name}"),
6095 }
6096 });
6097 });
6098 }
6099
6100 #[rstest::rstest]
6101 #[case("on_order_filled")]
6102 #[case("on_order_event")]
6103 fn test_python_order_event_callback_exception_continues_dispatch(
6104 #[case] method_name: &str,
6105 #[values(false, true)] shutdown: bool,
6106 ) {
6107 Python::initialize();
6108 Python::attach(|py| {
6109 let tracker = assert_python_callback(py, method_name, true, shutdown, |strategy| {
6110 Component::start(strategy.inner_mut()).unwrap();
6111 Strategy::handle_order_event(
6112 strategy.inner_mut(),
6113 OrderEventAny::Filled(OrderFilledSpec::builder().build()),
6114 );
6115 assert_eq!(Component::state(strategy.inner()), ComponentState::Running);
6116 Ok(())
6117 });
6118
6119 assert_eq!(
6120 python_method_call_names(&tracker, py),
6121 ["on_start", "on_order_filled", "on_order_event"],
6122 );
6123 });
6124 }
6125
6126 #[rstest::rstest]
6127 fn test_python_handle_order_event_dispatches_specific_and_aggregate_callbacks() {
6128 pyo3::Python::initialize();
6129 Python::attach(|py| {
6130 let (py_strategy, mut rust_strategy) = create_registered_tracking_strategy(py);
6131
6132 rust_strategy.py_start().unwrap();
6133 Strategy::handle_order_event(
6134 rust_strategy.inner_mut(),
6135 OrderEventAny::Accepted(OrderAccepted::default()),
6136 );
6137
6138 assert_eq!(
6139 python_method_call_count(&py_strategy, py, "on_order_accepted"),
6140 1
6141 );
6142 assert_eq!(
6143 python_method_call_count(&py_strategy, py, "on_order_event"),
6144 1
6145 );
6146 let call_names = python_method_call_names(&py_strategy, py);
6147 assert_eq!(
6148 &call_names[call_names.len() - 2..],
6149 ["on_order_accepted", "on_order_event"],
6150 );
6151 });
6152 }
6153
6154 #[rstest::rstest]
6155 #[case("on_position_event")]
6156 #[case("on_position_opened")]
6157 #[case("on_position_changed")]
6158 #[case("on_position_closed")]
6159 fn test_python_dispatch_position_callback_matrix(
6160 #[case] method_name: &str,
6161 #[values(false, true)] raises: bool,
6162 #[values(false, true)] shutdown: bool,
6163 ) {
6164 pyo3::Python::initialize();
6165 Python::attach(|py| {
6166 assert_python_callback(py, method_name, raises, shutdown, |rust_strategy| {
6167 match method_name {
6168 "on_position_event" => {
6169 Strategy::on_position_event(
6170 rust_strategy.inner_mut(),
6171 PositionEvent::PositionOpened(sample_position_opened()),
6172 );
6173 Ok(())
6174 }
6175 "on_position_opened" => {
6176 Strategy::on_position_opened(
6177 rust_strategy.inner_mut(),
6178 sample_position_opened(),
6179 );
6180 Ok(())
6181 }
6182 "on_position_changed" => {
6183 Strategy::on_position_changed(
6184 rust_strategy.inner_mut(),
6185 sample_position_changed(),
6186 );
6187 Ok(())
6188 }
6189 "on_position_closed" => {
6190 Strategy::on_position_closed(
6191 rust_strategy.inner_mut(),
6192 sample_position_closed(),
6193 );
6194 Ok(())
6195 }
6196 _ => unreachable!("unhandled position callback case: {method_name}"),
6197 }
6198 });
6199 });
6200 }
6201
6202 #[rstest::rstest]
6203 fn test_python_handle_position_event_dispatches_specific_and_aggregate_callbacks() {
6204 pyo3::Python::initialize();
6205 Python::attach(|py| {
6206 let (py_strategy, mut rust_strategy) = create_registered_tracking_strategy(py);
6207
6208 rust_strategy.py_start().unwrap();
6209 Strategy::handle_position_event(
6210 rust_strategy.inner_mut(),
6211 PositionEvent::PositionOpened(sample_position_opened()),
6212 );
6213
6214 assert_eq!(
6215 python_method_call_count(&py_strategy, py, "on_position_opened"),
6216 1
6217 );
6218 assert_eq!(
6219 python_method_call_count(&py_strategy, py, "on_position_event"),
6220 1
6221 );
6222 let call_names = python_method_call_names(&py_strategy, py);
6223 assert_eq!(
6224 &call_names[call_names.len() - 2..],
6225 ["on_position_opened", "on_position_event"],
6226 );
6227 });
6228 }
6229
6230 #[rstest::rstest]
6231 fn test_python_self_is_weak() {
6232 Python::initialize();
6233
6234 Python::attach(|py| {
6235 let instance = py
6236 .get_type::<PyStrategy>()
6237 .call0()
6238 .expect("Strategy should construct");
6239 let weakref =
6240 PyWeakrefReference::new(&instance).expect("Strategy should be weak-referenceable");
6241 assert!(weakref.upgrade().is_some());
6242
6243 drop(instance);
6244
6245 assert!(
6247 weakref.upgrade().is_none(),
6248 "an unregistered Strategy must be collected once its last Python owner is dropped",
6249 );
6250 });
6251 }
6252}