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