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