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