1use std::{
19 sync::{
20 Arc,
21 atomic::{AtomicBool, Ordering},
22 },
23 time::Duration,
24};
25
26use ahash::{AHashMap, AHashSet};
27use async_trait::async_trait;
28use nautilus_common::{
29 clients::DataClient,
30 live::runner::get_data_event_sender,
31 messages::{
32 DataEvent,
33 data::{
34 SubscribeBookDeltas, SubscribeInstrumentClose, SubscribeInstrumentStatus,
35 SubscribeTrades, UnsubscribeBars, UnsubscribeBookDeltas, UnsubscribeCustomData,
36 UnsubscribeInstrument, UnsubscribeInstrumentClose, UnsubscribeInstrumentStatus,
37 UnsubscribeInstruments, UnsubscribeQuotes, UnsubscribeTrades,
38 },
39 },
40 providers::InstrumentProvider,
41};
42use nautilus_core::{
43 AtomicMap, Params,
44 time::{AtomicTime, get_atomic_clock_realtime},
45};
46use nautilus_live::{
47 SocketControl,
48 task::{TaskGroup, TaskGroupGuard},
49};
50use nautilus_model::{
51 data::{CustomData, CustomDataTrait, Data, DataType, OrderBookDeltas, TradeTick},
52 identifiers::{ClientId, InstrumentId, TradeId, Venue},
53 instruments::{Instrument, InstrumentAny},
54 types::{Currency, Money},
55};
56use parking_lot::Mutex;
57use rust_decimal::Decimal;
58
59use crate::{
60 common::{
61 consts::{BETFAIR_RACE_STREAM_HOST, BETFAIR_VENUE},
62 credential::BetfairCredential,
63 enums::{MarketDataFilterField, MarketStatus, SegmentType},
64 parse::{
65 extract_market_id, make_instrument_id, parse_betfair_price, parse_betfair_quantity,
66 parse_market_definition, parse_millis_timestamp,
67 },
68 },
69 config::BetfairDataClientConfig,
70 data_types::{BetfairSequenceCompleted, register_betfair_custom_data},
71 http::client::BetfairHttpClient,
72 provider::{BetfairInstrumentProvider, NavigationFilter},
73 stream::{
74 CRICKET_STREAMS_ENDPOINT, DATA_STREAMS_ENDPOINT, RACE_STREAMS_ENDPOINT,
75 client::{
76 BetfairRaceStreamClient, BetfairStreamClient, HeartbeatTimeoutSource,
77 StreamMessageHandler,
78 },
79 config::BetfairStreamConfig,
80 messages::{MarketDataFilter, StreamMarketFilter, StreamMessage},
81 parse::{
82 make_trade_tick, parse_betfair_starting_prices, parse_betfair_ticker,
83 parse_bsp_book_deltas, parse_cricket_match, parse_instrument_closes,
84 parse_instrument_statuses, parse_race_progress, parse_race_runner_data,
85 parse_runner_book_deltas,
86 },
87 },
88};
89
90const KEEP_ALIVE_INTERVAL_SECS: u64 = 36_000;
92
93#[derive(Debug)]
95pub struct BetfairDataClient {
96 clock: &'static AtomicTime,
97 client_id: ClientId,
98 http_client: Arc<BetfairHttpClient>,
99 provider: BetfairInstrumentProvider,
100 stream_client: Option<Arc<BetfairStreamClient>>,
101 socket_control: Option<SocketControl>,
102 race_socket_control: Option<Arc<SocketControl>>,
103 race_stream_client: Option<Arc<BetfairRaceStreamClient>>,
104 cricket_socket_control: Option<Arc<SocketControl>>,
105 cricket_stream_client: Option<Arc<BetfairRaceStreamClient>>,
106 credential: BetfairCredential,
107 stream_config: BetfairStreamConfig,
108 config: BetfairDataClientConfig,
109 currency: Currency,
110 is_connected: AtomicBool,
111 data_sender: tokio::sync::mpsc::UnboundedSender<DataEvent>,
112 instruments: Arc<AtomicMap<InstrumentId, InstrumentAny>>,
113 subscribed_market_ids: AHashSet<String>,
114 session_tasks: TaskGroup,
115 command_tasks: TaskGroup,
116 stream_shutdowns: Arc<Mutex<Vec<BetfairStreamShutdown>>>,
117 shutdown_errors: Vec<String>,
118}
119
120pub(crate) fn custom_data_with_instrument(
123 value: Arc<dyn CustomDataTrait>,
124 instrument_id: InstrumentId,
125) -> CustomData {
126 let mut metadata = Params::new();
127 metadata.insert(
128 "instrument_id".to_string(),
129 serde_json::Value::String(instrument_id.to_string()),
130 );
131 let data_type = DataType::new(
132 value.type_name(),
133 Some(metadata),
134 Some(instrument_id.to_string()),
135 );
136 CustomData::new(value, data_type)
137}
138
139impl BetfairDataClient {
140 #[must_use]
142 #[expect(clippy::too_many_arguments)]
143 pub fn new(
144 client_id: ClientId,
145 http_client: BetfairHttpClient,
146 credential: BetfairCredential,
147 stream_config: BetfairStreamConfig,
148 config: BetfairDataClientConfig,
149 nav_filter: NavigationFilter,
150 currency: Currency,
151 min_notional: Option<Money>,
152 ) -> Self {
153 let data_sender = get_data_event_sender();
154 let http_client = Arc::new(http_client);
155 let socket_control = Some(SocketControl::new(
156 client_id,
157 Some(*BETFAIR_VENUE),
158 DATA_STREAMS_ENDPOINT,
159 ));
160 let race_socket_control = config.subscribe_race_data.then(|| {
161 Arc::new(SocketControl::new(
162 client_id,
163 Some(*BETFAIR_VENUE),
164 RACE_STREAMS_ENDPOINT,
165 ))
166 });
167 let cricket_socket_control = config.subscribe_cricket_data.then(|| {
168 Arc::new(SocketControl::new(
169 client_id,
170 Some(*BETFAIR_VENUE),
171 CRICKET_STREAMS_ENDPOINT,
172 ))
173 });
174 let provider = BetfairInstrumentProvider::new(
175 Arc::clone(&http_client),
176 nav_filter,
177 currency,
178 min_notional,
179 );
180
181 let session_tasks = TaskGroup::new();
182 let command_tasks = TaskGroup::new();
183
184 Self {
185 clock: get_atomic_clock_realtime(),
186 client_id,
187 http_client,
188 provider,
189 stream_client: None,
190 socket_control,
191 race_socket_control,
192 race_stream_client: None,
193 cricket_socket_control,
194 cricket_stream_client: None,
195 credential,
196 stream_config,
197 config,
198 currency,
199 is_connected: AtomicBool::new(false),
200 data_sender,
201 instruments: Arc::new(AtomicMap::new()),
202 subscribed_market_ids: AHashSet::new(),
203 session_tasks,
204 command_tasks,
205 stream_shutdowns: Arc::new(Mutex::new(Vec::new())),
206 shutdown_errors: Vec::new(),
207 }
208 }
209
210 fn spawn_command<F>(&self, future: F)
211 where
212 F: std::future::Future<Output = ()> + Send + 'static,
213 {
214 if let Err(e) = self.command_tasks.spawn(future) {
215 log::warn!("Skipping Betfair data command after shutdown began: {e}");
216 }
217 }
218
219 async fn finish_tasks(&self) -> anyhow::Result<()> {
220 let (session_result, command_result) = tokio::join!(
221 self.session_tasks
222 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
223 self.command_tasks
224 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
225 );
226 session_result
227 .map_err(|e| anyhow::anyhow!("Failed to finish Betfair data session tasks: {e}"))?;
228 command_result
229 .map_err(|e| anyhow::anyhow!("Failed to finish Betfair data command tasks: {e}"))?;
230 Ok(())
231 }
232
233 async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
234 if !self.session_tasks.is_open() || !self.command_tasks.is_open() {
235 self.teardown_partial_connect().await?;
236 self.session_tasks
237 .start_generation()
238 .map_err(|e| anyhow::anyhow!("Failed to start Betfair data session tasks: {e}"))?;
239 self.command_tasks
240 .start_generation()
241 .map_err(|e| anyhow::anyhow!("Failed to start Betfair data command tasks: {e}"))?;
242 }
243 Ok(())
244 }
245
246 fn begin_stream_shutdown(&self) {
247 for stream in self.stream_shutdowns.lock().iter() {
248 stream.begin_shutdown();
249 }
250 }
251
252 async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
253 self.session_tasks.begin_shutdown();
254 self.command_tasks.begin_shutdown();
255 self.begin_stream_shutdown();
256 self.is_connected.store(false, Ordering::Relaxed);
257
258 if let Some(client) = self.cricket_stream_client.as_ref() {
259 client.close().await;
260 self.cricket_stream_client = None;
261 }
262
263 if let Some(client) = self.race_stream_client.as_ref() {
264 client.close().await;
265 self.race_stream_client = None;
266 }
267
268 if let Some(client) = self.stream_client.as_ref() {
269 match client.close().await {
270 Ok(()) => self.stream_client = None,
271 Err(e) => self
272 .shutdown_errors
273 .push(format!("stream shutdown failed: {e}")),
274 }
275 }
276
277 self.http_client.disconnect().await;
278
279 if let Err(e) = self.finish_tasks().await {
280 self.shutdown_errors.push(e.to_string());
281 }
282 self.is_connected.store(false, Ordering::Release);
283 self.deregister_socket_controls();
284
285 if self.stream_client.is_none()
286 && self.race_stream_client.is_none()
287 && self.cricket_stream_client.is_none()
288 {
289 self.stream_shutdowns.lock().clear();
290 }
291
292 if self.shutdown_errors.is_empty() {
293 Ok(())
294 } else {
295 let errors = std::mem::take(&mut self.shutdown_errors);
296 anyhow::bail!("Betfair data shutdown failed: {}", errors.join("; "))
297 }
298 }
299
300 fn create_stream_handler(
301 data_sender: tokio::sync::mpsc::UnboundedSender<DataEvent>,
302 instruments: Arc<AtomicMap<InstrumentId, InstrumentAny>>,
303 currency: Currency,
304 min_notional: Option<Money>,
305 reconnect_tx: tokio::sync::mpsc::UnboundedSender<()>,
306 clock: &'static AtomicTime,
307 ) -> StreamMessageHandler {
308 let traded_volumes: Arc<Mutex<AHashMap<(InstrumentId, Decimal), Decimal>>> =
311 Arc::new(Mutex::new(AHashMap::new()));
312 let has_initial_connection = Arc::new(AtomicBool::new(false));
313
314 Arc::new(move |msg: StreamMessage| {
315 let ts_init = clock.get_time_ns();
316
317 match msg {
318 StreamMessage::MarketChange(mcm) => {
319 if mcm.is_heartbeat() {
320 return;
321 }
322
323 let sequence_complete = mcm
324 .segment_type
325 .is_none_or(|segment| segment == SegmentType::SegEnd);
326
327 let Some(market_changes) = &mcm.mc else {
328 return;
329 };
330
331 let ts_event = parse_millis_timestamp(mcm.pt);
332
333 for mc in market_changes {
334 let is_snapshot = mc.img;
335 let mut market_closed = false;
336
337 if let Some(def) = &mc.market_definition {
338 match parse_market_definition(
342 &mc.id,
343 def,
344 currency,
345 ts_event,
346 ts_init,
347 min_notional,
348 ) {
349 Ok(new_instruments) => {
350 instruments.rcu(|m| {
351 for inst in &new_instruments {
352 m.insert(inst.id(), inst.clone());
353 }
354 });
355
356 for inst in new_instruments {
357 if let Err(e) =
358 data_sender.send(DataEvent::Instrument(inst))
359 {
360 log::warn!("Failed to send instrument: {e}");
361 }
362 }
363 }
364 Err(e) => {
365 log::warn!(
366 "Failed to parse market definition for {}: {e}",
367 mc.id
368 );
369 }
370 }
371
372 if let Some(status) = &def.status {
373 market_closed = *status == MarketStatus::Closed;
374
375 for event in
376 parse_instrument_statuses(&mc.id, def, ts_event, ts_init)
377 {
378 if let Err(e) =
379 data_sender.send(DataEvent::InstrumentStatus(event))
380 {
381 log::warn!("Failed to send instrument status: {e}");
382 }
383 }
384 }
385
386 for sp in parse_betfair_starting_prices(&mc.id, def, ts_event, ts_init)
387 {
388 let instrument_id = sp.instrument_id;
389 let custom =
390 custom_data_with_instrument(Arc::new(sp), instrument_id);
391
392 if let Err(e) =
393 data_sender.send(DataEvent::Data(Data::Custom(custom)))
394 {
395 log::warn!("Failed to send starting price: {e}");
396 }
397 }
398
399 for close in parse_instrument_closes(&mc.id, def, ts_event, ts_init) {
400 if let Err(e) =
401 data_sender.send(DataEvent::Data(Data::InstrumentClose(close)))
402 {
403 log::warn!("Failed to send instrument close: {e}");
404 }
405 }
406 }
407
408 let mut buffered_deltas: Vec<OrderBookDeltas> = Vec::new();
412 let mut buffered_bsp_customs: Vec<CustomData> = Vec::new();
413
414 if let Some(runner_changes) = &mc.rc {
415 for rc in runner_changes {
416 let handicap = rc.hc.unwrap_or(Decimal::ZERO);
417 let instrument_id = make_instrument_id(&mc.id, rc.id, handicap);
418
419 match parse_runner_book_deltas(
420 instrument_id,
421 rc,
422 is_snapshot,
423 mcm.pt,
424 ts_event,
425 ts_init,
426 ) {
427 Ok(Some(deltas)) => {
428 if is_snapshot {
429 if let Err(e) = data_sender.send(DataEvent::Data(
430 Data::Deltas(Box::new(deltas)),
431 )) {
432 log::warn!("Failed to send book deltas: {e}");
433 }
434 } else {
435 buffered_deltas.push(deltas);
436 }
437 }
438 Ok(None) => {}
439 Err(e) => {
440 log::warn!(
441 "Failed to parse book deltas for {instrument_id}: {e}"
442 );
443 }
444 }
445
446 if let Some(trades) = &rc.trd {
447 let mut volumes = traded_volumes.lock();
448
449 for pv in trades {
450 if pv.volume == Decimal::ZERO {
451 continue;
452 }
453
454 let key = (instrument_id, pv.price);
455 let prev_volume =
456 volumes.get(&key).copied().unwrap_or(Decimal::ZERO);
457
458 if pv.volume <= prev_volume {
459 volumes.insert(key, pv.volume);
460 continue;
461 }
462
463 let trade_volume = pv.volume - prev_volume;
464 volumes.insert(key, pv.volume);
465
466 let price = match parse_betfair_price(pv.price) {
467 Ok(p) => p,
468 Err(e) => {
469 log::warn!("Invalid trade price: {e}");
470 continue;
471 }
472 };
473 let size = match parse_betfair_quantity(trade_volume) {
474 Ok(q) => q,
475 Err(e) => {
476 log::warn!("Invalid trade size: {e}");
477 continue;
478 }
479 };
480 let trade_id = TradeId::new(format!(
481 "{}-{}-{}",
482 mcm.pt, rc.id, pv.price
483 ));
484 let tick: TradeTick = make_trade_tick(
485 instrument_id,
486 price,
487 size,
488 trade_id,
489 ts_event,
490 ts_init,
491 );
492
493 if let Err(e) =
494 data_sender.send(DataEvent::Data(Data::Trade(tick)))
495 {
496 log::warn!("Failed to send trade tick: {e}");
497 }
498 }
499 }
500
501 if let Some(ticker) =
502 parse_betfair_ticker(instrument_id, rc, ts_event, ts_init)
503 {
504 let custom = custom_data_with_instrument(
505 Arc::new(ticker),
506 instrument_id,
507 );
508
509 if let Err(e) =
510 data_sender.send(DataEvent::Data(Data::Custom(custom)))
511 {
512 log::warn!("Failed to send ticker: {e}");
513 }
514 }
515
516 for bsp_delta in
517 parse_bsp_book_deltas(instrument_id, rc, ts_event, ts_init)
518 {
519 buffered_bsp_customs.push(custom_data_with_instrument(
520 Arc::new(bsp_delta),
521 instrument_id,
522 ));
523 }
524 }
525 }
526
527 for deltas in buffered_deltas {
528 if let Err(e) =
529 data_sender.send(DataEvent::Data(Data::Deltas(Box::new(deltas))))
530 {
531 log::warn!("Failed to send book deltas: {e}");
532 }
533 }
534
535 for custom in buffered_bsp_customs {
536 if let Err(e) = data_sender.send(DataEvent::Data(Data::Custom(custom)))
537 {
538 log::warn!("Failed to send BSP book delta: {e}");
539 }
540 }
541
542 if market_closed {
543 let prefix = format!("{}-", mc.id);
544
545 traded_volumes
546 .lock()
547 .retain(|k, _| !k.0.symbol.as_str().starts_with(&prefix));
548 }
549 }
550
551 if sequence_complete {
552 let completed = BetfairSequenceCompleted::new(ts_event, ts_init);
553 let custom = CustomData::from_arc(Arc::new(completed));
554 if let Err(e) = data_sender.send(DataEvent::Data(Data::Custom(custom))) {
555 log::warn!("Failed to send sequence completed: {e}");
556 }
557 }
558 }
559 StreamMessage::Connection(_) => {
560 if has_initial_connection.swap(true, Ordering::SeqCst) {
561 log::info!("Betfair data stream reconnected");
562 let _ = reconnect_tx.send(());
563 } else {
564 log::debug!("Betfair data stream connected");
565 }
566 }
567 StreamMessage::Status(status) => {
568 if status.connection_closed {
569 log::warn!(
570 "Betfair stream closed: {:?} - {:?}",
571 status.error_code,
572 status.error_message,
573 );
574 }
575 }
576 StreamMessage::RaceChange(rcm) => {
577 if let Some(race_changes) = &rcm.rc {
578 let ts_event_fallback = parse_millis_timestamp(rcm.pt);
579
580 for rc in race_changes {
581 let race_id = rc.id.as_deref().unwrap_or("");
582 let market_id = rc.mid.as_deref().unwrap_or("");
583
584 if let Some(runners) = &rc.rrc {
585 for rrc in runners {
586 let ts_event =
587 rrc.ft.map_or(ts_event_fallback, parse_millis_timestamp);
588
589 if let Some(runner) = parse_race_runner_data(
590 race_id, market_id, rrc, ts_event, ts_init,
591 ) {
592 let selection_id = rrc.id.unwrap_or(0);
593 let mut metadata = Params::new();
594 metadata.insert(
595 "selection_id".to_string(),
596 serde_json::Value::Number(selection_id.into()),
597 );
598 let value: Arc<dyn CustomDataTrait> = Arc::new(runner);
599 let data_type =
600 DataType::new(value.type_name(), Some(metadata), None);
601 let custom = CustomData::new(value, data_type);
602
603 if let Err(e) =
604 data_sender.send(DataEvent::Data(Data::Custom(custom)))
605 {
606 log::warn!("Failed to send race runner data: {e}");
607 }
608 }
609 }
610 }
611
612 if let Some(rpc) = &rc.rpc {
613 let ts_event =
614 rpc.ft.map_or(ts_event_fallback, parse_millis_timestamp);
615
616 let progress =
617 parse_race_progress(race_id, market_id, rpc, ts_event, ts_init);
618 let mut metadata = Params::new();
619 metadata.insert(
620 "race_id".to_string(),
621 serde_json::Value::String(race_id.to_string()),
622 );
623 let value: Arc<dyn CustomDataTrait> = Arc::new(progress);
624 let data_type =
625 DataType::new(value.type_name(), Some(metadata), None);
626 let custom = CustomData::new(value, data_type);
627
628 if let Err(e) =
629 data_sender.send(DataEvent::Data(Data::Custom(custom)))
630 {
631 log::warn!("Failed to send race progress: {e}");
632 }
633 }
634 }
635 }
636 }
637 StreamMessage::CricketChange(ccm) => {
638 if let Some(cricket_changes) = &ccm.cc {
639 let ts_event = parse_millis_timestamp(ccm.pt);
640
641 for cricket_change in cricket_changes {
642 if let Some(cricket) =
643 parse_cricket_match(cricket_change, ts_event, ts_init)
644 {
645 let mut metadata = Params::new();
646 metadata.insert(
647 "event_id".to_string(),
648 serde_json::Value::String(cricket.event_id.clone()),
649 );
650 let value: Arc<dyn CustomDataTrait> = Arc::new(cricket);
651 let data_type =
652 DataType::new(value.type_name(), Some(metadata), None);
653 let custom = CustomData::new(value, data_type);
654
655 if let Err(e) =
656 data_sender.send(DataEvent::Data(Data::Custom(custom)))
657 {
658 log::warn!("Failed to send cricket match: {e}");
659 }
660 }
661 }
662 }
663 }
664 StreamMessage::OrderChange(_) => {}
665 }
666 })
667 }
668}
669
670#[async_trait(?Send)]
671impl DataClient for BetfairDataClient {
672 fn client_id(&self) -> ClientId {
673 self.client_id
674 }
675
676 fn venue(&self) -> Option<Venue> {
677 Some(*BETFAIR_VENUE)
678 }
679
680 fn start(&mut self) -> anyhow::Result<()> {
681 log::info!("Starting Betfair data client: {}", self.client_id);
682 Ok(())
683 }
684
685 fn stop(&mut self) -> anyhow::Result<()> {
686 log::info!("Stopping Betfair data client: {}", self.client_id);
687
688 self.session_tasks.begin_shutdown();
689 self.command_tasks.begin_shutdown();
690 self.begin_stream_shutdown();
691 self.is_connected.store(false, Ordering::Relaxed);
692
693 Ok(())
694 }
695
696 fn reset(&mut self) -> anyhow::Result<()> {
697 log::info!("Resetting Betfair data client: {}", self.client_id);
698
699 self.session_tasks.begin_shutdown();
700 self.command_tasks.begin_shutdown();
701 self.begin_stream_shutdown();
702
703 self.provider.store_mut().clear();
704 self.subscribed_market_ids.clear();
705
706 self.instruments.store(AHashMap::new());
707 Ok(())
708 }
709
710 fn dispose(&mut self) -> anyhow::Result<()> {
711 log::debug!("Disposing Betfair data client: {}", self.client_id);
712 self.stop()
713 }
714
715 fn is_connected(&self) -> bool {
716 self.is_connected.load(Ordering::SeqCst)
717 && self.stream_client.as_ref().is_some_and(|client| {
718 client.is_authenticated()
719 && (self.subscribed_market_ids.is_empty() || client.is_market_ready())
720 })
721 }
722
723 fn is_disconnected(&self) -> bool {
724 !self.is_connected()
725 }
726
727 async fn connect(&mut self) -> anyhow::Result<()> {
728 if self.is_connected.load(Ordering::Acquire)
729 && self.session_tasks.is_open()
730 && self.command_tasks.is_open()
731 {
732 return Ok(());
733 }
734
735 self.prepare_task_groups().await?;
736 let stream_shutdowns = Arc::clone(&self.stream_shutdowns);
737 let setup_guard =
738 TaskGroupGuard::new(&[&self.session_tasks, &self.command_tasks], move || {
739 for stream in stream_shutdowns.lock().iter() {
740 stream.begin_shutdown();
741 }
742 });
743
744 register_betfair_custom_data();
745
746 self.http_client
747 .connect()
748 .await
749 .map_err(|e| anyhow::anyhow!("{e}"))?;
750
751 self.provider.load_all(None).await?;
752
753 let loaded: Vec<InstrumentAny> = self
754 .provider
755 .store()
756 .list_all()
757 .into_iter()
758 .cloned()
759 .collect();
760
761 self.instruments.rcu(|m| {
762 for inst in &loaded {
763 m.insert(inst.id(), inst.clone());
764 }
765 });
766
767 for inst in &loaded {
768 if let Err(e) = self.data_sender.send(DataEvent::Instrument(inst.clone())) {
769 log::warn!("Failed to send instrument: {e}");
770 }
771 }
772
773 log::debug!("Cached {} instruments for {}", loaded.len(), self.client_id,);
774
775 let session_token = self
776 .http_client
777 .session_token()
778 .await
779 .ok_or_else(|| anyhow::anyhow!("No session token after login"))?;
780
781 let (reconnect_tx, mut reconnect_rx) = tokio::sync::mpsc::unbounded_channel();
782
783 let handler = Self::create_stream_handler(
784 self.data_sender.clone(),
785 Arc::clone(&self.instruments),
786 self.currency,
787 self.provider.min_notional(),
788 reconnect_tx.clone(),
789 self.clock,
790 );
791
792 let state_sink = self.socket_control.as_ref().map(SocketControl::sink);
793 let stream_client = BetfairStreamClient::connect_with_state_sink(
794 &self.credential,
795 session_token,
796 handler,
797 self.stream_config.clone(),
798 HeartbeatTimeoutSource::Server,
799 state_sink,
800 )
801 .await
802 .map_err(|e| anyhow::anyhow!("{e}"))?;
803
804 let stream_client = Arc::new(stream_client);
805 if let Some(control) = &self.socket_control {
806 let reconnect_stream = Arc::clone(&stream_client);
807 control.register(move || reconnect_stream.request_reconnect_outcome());
808 }
809 self.stream_client = Some(stream_client);
810 self.stream_shutdowns
811 .lock()
812 .push(BetfairStreamShutdown::Exchange(Arc::clone(
813 self.stream_client.as_ref().expect("stream client assigned"),
814 )));
815
816 let session_result = async {
817 if self.config.subscribe_race_data {
818 let race_config = BetfairStreamConfig {
819 host: BETFAIR_RACE_STREAM_HOST.to_string(),
820 ..self.stream_config.clone()
821 };
822
823 let race_session = self
824 .http_client
825 .session_token()
826 .await
827 .ok_or_else(|| anyhow::anyhow!("No session token for race stream"))?;
828
829 let race_handler = Self::create_stream_handler(
830 self.data_sender.clone(),
831 Arc::clone(&self.instruments),
832 self.currency,
833 self.provider.min_notional(),
834 reconnect_tx.clone(),
835 self.clock,
836 );
837
838 let (race_fatal_tx, mut race_fatal_rx) = tokio::sync::mpsc::unbounded_channel();
839
840 let state_sink = self
841 .race_socket_control
842 .as_ref()
843 .map(|control| control.sink());
844
845 match BetfairRaceStreamClient::connect_decoded(
846 &self.credential,
847 race_session,
848 race_handler,
849 race_config,
850 race_fatal_tx,
851 state_sink,
852 )
853 .await
854 {
855 Ok(client) => {
856 let race_client = Arc::new(client);
857 if let Some(control) = &self.race_socket_control {
858 let reconnect_client = Arc::clone(&race_client);
859 control.register(move || reconnect_client.request_reconnect_outcome());
860 }
861 self.race_stream_client = Some(Arc::clone(&race_client));
862 self.stream_shutdowns
863 .lock()
864 .push(BetfairStreamShutdown::Auxiliary(Arc::clone(&race_client)));
865
866 let race_socket_control = self.race_socket_control.as_ref().map(Arc::clone);
867
868 self.session_tasks
869 .spawn(async move {
870 if race_fatal_rx.recv().await.is_some() {
871 log::error!(
872 "Betfair race stream permanently disabled due to fatal error"
873 );
874 race_client.close().await;
875
876 if let Some(control) = race_socket_control {
877 control.deregister();
878 }
879 }
880 })
881 .map_err(|e| {
882 anyhow::anyhow!("Failed to register Betfair race fatal task: {e}")
883 })?;
884
885 log::debug!("Betfair race stream connected");
886 }
887 Err(e) => {
888 log::warn!("Betfair race stream connect failed: {e}");
889
890 if let Some(control) = &self.race_socket_control {
891 control.deregister();
892 }
893 self.race_stream_client = None;
894 }
895 }
896 }
897
898 if self.config.subscribe_cricket_data {
899 let cricket_config = BetfairStreamConfig {
900 host: BETFAIR_RACE_STREAM_HOST.to_string(),
901 ..self.stream_config.clone()
902 };
903
904 let cricket_session = self
905 .http_client
906 .session_token()
907 .await
908 .ok_or_else(|| anyhow::anyhow!("No session token for cricket stream"))?;
909
910 let cricket_handler = Self::create_stream_handler(
911 self.data_sender.clone(),
912 Arc::clone(&self.instruments),
913 self.currency,
914 self.provider.min_notional(),
915 reconnect_tx.clone(),
916 self.clock,
917 );
918
919 let (cricket_fatal_tx, mut cricket_fatal_rx) =
920 tokio::sync::mpsc::unbounded_channel();
921
922 let state_sink = self
923 .cricket_socket_control
924 .as_ref()
925 .map(|control| control.sink());
926
927 match BetfairRaceStreamClient::connect_cricket_decoded(
928 &self.credential,
929 cricket_session,
930 cricket_handler,
931 cricket_config,
932 cricket_fatal_tx,
933 state_sink,
934 )
935 .await
936 {
937 Ok(client) => {
938 let cricket_client = Arc::new(client);
939 if let Some(control) = &self.cricket_socket_control {
940 let reconnect_client = Arc::clone(&cricket_client);
941 control.register(move || reconnect_client.request_reconnect_outcome());
942 }
943 self.cricket_stream_client = Some(Arc::clone(&cricket_client));
944 self.stream_shutdowns
945 .lock()
946 .push(BetfairStreamShutdown::Auxiliary(Arc::clone(
947 &cricket_client,
948 )));
949
950 let cricket_socket_control =
951 self.cricket_socket_control.as_ref().map(Arc::clone);
952
953 self.session_tasks
954 .spawn(async move {
955 if cricket_fatal_rx.recv().await.is_some() {
956 log::error!(
957 "Betfair cricket stream permanently disabled due to fatal error"
958 );
959 cricket_client.close().await;
960
961 if let Some(control) = cricket_socket_control {
962 control.deregister();
963 }
964 }
965 })
966 .map_err(|e| {
967 anyhow::anyhow!("Failed to register Betfair cricket fatal task: {e}")
968 })?;
969
970 log::debug!("Betfair cricket stream connected");
971 }
972 Err(e) => {
973 log::warn!("Betfair cricket stream connect failed: {e}");
974
975 if let Some(control) = &self.cricket_socket_control {
976 control.deregister();
977 }
978 self.cricket_stream_client = None;
979 }
980 }
981 }
982
983 let keep_alive_client = Arc::clone(&self.http_client);
984 let keep_alive_stream = Arc::clone(self.stream_client.as_ref().unwrap());
985 let keep_alive_race_stream = self.race_stream_client.as_ref().map(Arc::clone);
986 let keep_alive_cricket_stream = self.cricket_stream_client.as_ref().map(Arc::clone);
987 let keep_alive_app_key = self.credential.app_key().to_string();
988
989 self.session_tasks
990 .spawn(async move {
991 let interval = tokio::time::Duration::from_secs(KEEP_ALIVE_INTERVAL_SECS);
992 loop {
993 tokio::time::sleep(interval).await;
994
995 let session_replaced = match keep_alive_client.keep_alive_with_token().await
996 {
997 Ok(_) => false,
998 Err(ref e) if e.is_login_failed() => {
999 log::warn!("Betfair session expired, attempting re-login: {e}");
1000
1001 match keep_alive_client.reconnect_with_token().await {
1002 Ok(_) => true,
1003 Err(e) => {
1004 log::warn!("Betfair re-login failed: {e}");
1005 continue;
1006 }
1007 }
1008 }
1009 Err(e) => {
1010 log::warn!("Betfair keep-alive failed (transient): {e}");
1011 continue;
1012 }
1013 };
1014
1015 let _ = keep_alive_client
1016 .with_session_token(|token| {
1017 refresh_stream_sessions(
1018 keep_alive_stream.as_ref(),
1019 keep_alive_race_stream.as_deref(),
1020 keep_alive_cricket_stream.as_deref(),
1021 &keep_alive_app_key,
1022 token,
1023 session_replaced,
1024 );
1025 })
1026 .await;
1027 log::debug!("Betfair session keep-alive sent");
1028 }
1029 })
1030 .map_err(|e| anyhow::anyhow!("Failed to register Betfair keep-alive task: {e}"))?;
1031
1032 let reconnect_http = Arc::clone(&self.http_client);
1033 let reconnect_stream = Arc::clone(self.stream_client.as_ref().unwrap());
1034 let reconnect_race_stream = self.race_stream_client.as_ref().map(Arc::clone);
1035 let reconnect_cricket_stream = self.cricket_stream_client.as_ref().map(Arc::clone);
1036 let reconnect_app_key = self.credential.app_key().to_string();
1037
1038 self.session_tasks
1039 .spawn(async move {
1040 while reconnect_rx.recv().await.is_some() {
1041 log::info!("Handling data stream reconnection");
1042
1043 let session_replaced = match reconnect_http.keep_alive_with_token().await {
1044 Ok(_) => false,
1045 Err(ref e) if e.is_login_failed() => {
1046 log::warn!(
1047 "Session expired on reconnect, attempting re-login: {e}"
1048 );
1049
1050 match reconnect_http.reconnect_with_token().await {
1051 Ok(_) => true,
1052 Err(e) => {
1053 log::warn!("Re-login failed on reconnect: {e}");
1054 continue;
1055 }
1056 }
1057 }
1058 Err(e) => {
1059 log::warn!("Keep-alive failed on reconnect (transient): {e}");
1060 continue;
1061 }
1062 };
1063
1064 let _ = reconnect_http
1065 .with_session_token(|token| {
1066 refresh_stream_sessions(
1067 reconnect_stream.as_ref(),
1068 reconnect_race_stream.as_deref(),
1069 reconnect_cricket_stream.as_deref(),
1070 &reconnect_app_key,
1071 token,
1072 session_replaced,
1073 );
1074 })
1075 .await;
1076 }
1077 })
1078 .map_err(|e| anyhow::anyhow!("Failed to register Betfair reconnect task: {e}"))?;
1079
1080 Ok::<(), anyhow::Error>(())
1081 }
1082 .await;
1083
1084 if let Err(e) = session_result {
1085 if let Err(teardown_error) = self.teardown_partial_connect().await {
1086 return Err(e.context(format!(
1087 "Betfair data startup teardown failed: {teardown_error}"
1088 )));
1089 }
1090 return Err(e);
1091 }
1092
1093 self.is_connected.store(true, Ordering::Release);
1094 setup_guard.disarm();
1095
1096 log::info!("Betfair data client connected: {}", self.client_id);
1097 Ok(())
1098 }
1099
1100 async fn disconnect(&mut self) -> anyhow::Result<()> {
1101 self.teardown_partial_connect().await?;
1102 self.subscribed_market_ids.clear();
1103
1104 log::info!("Betfair data client disconnected: {}", self.client_id);
1105 Ok(())
1106 }
1107
1108 fn subscribe_book_deltas(&mut self, cmd: SubscribeBookDeltas) -> anyhow::Result<()> {
1109 let instrument_id = cmd.instrument_id;
1110 let market_id = extract_market_id(&instrument_id)?;
1111
1112 if !self.subscribed_market_ids.insert(market_id.clone()) {
1113 log::debug!("Book deltas already subscribed for market {market_id}");
1114 return Ok(());
1115 }
1116
1117 let stream_client = Arc::clone(
1118 self.stream_client
1119 .as_ref()
1120 .ok_or_else(|| anyhow::anyhow!("Stream client not connected"))?,
1121 );
1122
1123 let all_ids: Vec<String> = self.subscribed_market_ids.iter().cloned().collect();
1124
1125 let market_filter = StreamMarketFilter {
1126 market_ids: Some(all_ids),
1127 ..Default::default()
1128 };
1129
1130 let data_filter = MarketDataFilter {
1131 fields: Some(vec![
1132 MarketDataFilterField::ExAllOffers,
1133 MarketDataFilterField::ExTraded,
1134 MarketDataFilterField::ExTradedVol,
1135 MarketDataFilterField::ExLtp,
1136 MarketDataFilterField::ExMarketDef,
1137 MarketDataFilterField::SpTraded,
1138 MarketDataFilterField::SpProjected,
1139 ]),
1140 ladder_levels: None,
1141 };
1142
1143 let conflate_ms = self.config.stream_conflate_ms;
1144
1145 self.spawn_command(async move {
1146 if let Err(e) = stream_client
1147 .subscribe_markets(market_filter, data_filter, None, conflate_ms)
1148 .await
1149 {
1150 log::warn!("Failed to subscribe to market data: {e}");
1151 }
1152 });
1153
1154 Ok(())
1155 }
1156
1157 fn unsubscribe_book_deltas(&mut self, cmd: &UnsubscribeBookDeltas) -> anyhow::Result<()> {
1158 log::debug!(
1159 "Skipping unsubscribe book deltas for Betfair: {}",
1160 cmd.instrument_id
1161 );
1162 Ok(())
1163 }
1164
1165 fn subscribe_trades(&mut self, cmd: SubscribeTrades) -> anyhow::Result<()> {
1166 log::debug!(
1168 "Trade data included in book subscription for {}",
1169 cmd.instrument_id
1170 );
1171 Ok(())
1172 }
1173
1174 fn unsubscribe_trades(&mut self, cmd: &UnsubscribeTrades) -> anyhow::Result<()> {
1175 log::debug!(
1176 "Skipping unsubscribe trades for Betfair: {}",
1177 cmd.instrument_id
1178 );
1179 Ok(())
1180 }
1181
1182 fn subscribe_instrument_status(
1183 &mut self,
1184 cmd: SubscribeInstrumentStatus,
1185 ) -> anyhow::Result<()> {
1186 log::debug!(
1188 "Instrument status included in book subscription for {}",
1189 cmd.instrument_id
1190 );
1191 Ok(())
1192 }
1193
1194 fn unsubscribe_instrument_status(
1195 &mut self,
1196 cmd: &UnsubscribeInstrumentStatus,
1197 ) -> anyhow::Result<()> {
1198 log::debug!(
1199 "Skipping unsubscribe instrument status for Betfair: {}",
1200 cmd.instrument_id
1201 );
1202 Ok(())
1203 }
1204
1205 fn subscribe_instrument_close(&mut self, cmd: SubscribeInstrumentClose) -> anyhow::Result<()> {
1206 log::debug!(
1209 "Instrument close included in book subscription for {}",
1210 cmd.instrument_id
1211 );
1212 Ok(())
1213 }
1214
1215 fn unsubscribe_instrument_close(
1216 &mut self,
1217 cmd: &UnsubscribeInstrumentClose,
1218 ) -> anyhow::Result<()> {
1219 log::debug!(
1220 "Skipping unsubscribe instrument close for Betfair: {}",
1221 cmd.instrument_id
1222 );
1223 Ok(())
1224 }
1225
1226 fn unsubscribe(&mut self, _cmd: &UnsubscribeCustomData) -> anyhow::Result<()> {
1227 log::debug!("Skipping unsubscribe custom data for Betfair");
1228 Ok(())
1229 }
1230
1231 fn unsubscribe_instrument(&mut self, cmd: &UnsubscribeInstrument) -> anyhow::Result<()> {
1232 log::debug!(
1233 "Skipping unsubscribe instrument for Betfair: {}",
1234 cmd.instrument_id
1235 );
1236 Ok(())
1237 }
1238
1239 fn unsubscribe_instruments(&mut self, _cmd: &UnsubscribeInstruments) -> anyhow::Result<()> {
1240 log::debug!("Skipping unsubscribe instruments for Betfair");
1241 Ok(())
1242 }
1243
1244 fn unsubscribe_quotes(&mut self, cmd: &UnsubscribeQuotes) -> anyhow::Result<()> {
1245 log::debug!(
1246 "Skipping unsubscribe quotes for Betfair: {}",
1247 cmd.instrument_id
1248 );
1249 Ok(())
1250 }
1251
1252 fn unsubscribe_bars(&mut self, cmd: &UnsubscribeBars) -> anyhow::Result<()> {
1253 log::debug!("Skipping unsubscribe bars for Betfair: {}", cmd.bar_type);
1254 Ok(())
1255 }
1256}
1257
1258impl BetfairDataClient {
1259 fn deregister_socket_controls(&self) {
1260 let controls = [
1261 self.socket_control.as_ref(),
1262 self.race_socket_control.as_deref(),
1263 self.cricket_socket_control.as_deref(),
1264 ];
1265
1266 for control in controls.into_iter().flatten() {
1267 control.deregister();
1268 }
1269 }
1270}
1271
1272#[derive(Clone, Debug)]
1273enum BetfairStreamShutdown {
1274 Exchange(Arc<BetfairStreamClient>),
1275 Auxiliary(Arc<BetfairRaceStreamClient>),
1276}
1277
1278impl BetfairStreamShutdown {
1279 fn begin_shutdown(&self) {
1280 match self {
1281 Self::Exchange(client) => client.begin_shutdown(),
1282 Self::Auxiliary(client) => client.begin_shutdown(),
1283 }
1284 }
1285}
1286
1287fn refresh_stream_sessions(
1288 stream: &BetfairStreamClient,
1289 race_stream: Option<&BetfairRaceStreamClient>,
1290 cricket_stream: Option<&BetfairRaceStreamClient>,
1291 app_key: &str,
1292 token: &str,
1293 session_replaced: bool,
1294) {
1295 stream.update_auth(app_key, token.to_string());
1296
1297 if let Some(race_stream) = race_stream {
1298 race_stream.update_auth(app_key, token.to_string());
1299 }
1300
1301 if let Some(cricket_stream) = cricket_stream {
1302 cricket_stream.update_auth(app_key, token.to_string());
1303 }
1304
1305 if !session_replaced {
1306 return;
1307 }
1308
1309 let _ = stream.request_reconnect();
1310
1311 if let Some(race_stream) = race_stream {
1312 let _ = race_stream.request_reconnect();
1313 }
1314
1315 if let Some(cricket_stream) = cricket_stream {
1316 let _ = cricket_stream.request_reconnect();
1317 }
1318}
1319
1320#[cfg(test)]
1321mod tests {
1322 use nautilus_core::UnixNanos;
1323 use rstest::rstest;
1324
1325 use super::*;
1326 use crate::{
1327 common::testing::load_test_json,
1328 data_types::{BetfairCricketMatch, BetfairRaceRunnerData, BetfairSequenceCompleted},
1329 stream::messages::stream_decode,
1330 };
1331
1332 fn stream_handler_at(
1333 ts_init: UnixNanos,
1334 ) -> (
1335 StreamMessageHandler,
1336 tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
1337 ) {
1338 let (data_tx, data_rx) = tokio::sync::mpsc::unbounded_channel();
1339 let (reconnect_tx, _reconnect_rx) = tokio::sync::mpsc::unbounded_channel();
1340 let clock = Box::leak(Box::new(AtomicTime::new(false, ts_init)));
1341 let handler = BetfairDataClient::create_stream_handler(
1342 data_tx,
1343 Arc::new(AtomicMap::new()),
1344 Currency::GBP(),
1345 None,
1346 reconnect_tx,
1347 clock,
1348 );
1349
1350 (handler, data_rx)
1351 }
1352
1353 fn receive_custom<T: 'static>(
1354 data_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
1355 ) -> Arc<dyn CustomDataTrait> {
1356 while let Ok(event) = data_rx.try_recv() {
1357 if let DataEvent::Data(Data::Custom(custom)) = event
1358 && custom.data.as_any().is::<T>()
1359 {
1360 return custom.data;
1361 }
1362 }
1363
1364 panic!("expected {} custom data", std::any::type_name::<T>());
1365 }
1366
1367 #[rstest]
1368 fn test_stream_handler_sets_mcm_init_from_clock() {
1369 let ts_init = UnixNanos::from(1_800_000_000_000_000_001);
1370
1371 let (handler, mut data_rx) = stream_handler_at(ts_init);
1372 let data = load_test_json("stream/mcm_UPDATE.json");
1373
1374 handler(stream_decode(data.as_bytes()).unwrap());
1375
1376 let custom = receive_custom::<BetfairSequenceCompleted>(&mut data_rx);
1377 let completed = custom
1378 .as_any()
1379 .downcast_ref::<BetfairSequenceCompleted>()
1380 .unwrap();
1381
1382 assert_eq!(
1383 completed.ts_event,
1384 UnixNanos::from(1_471_370_160_471_000_000)
1385 );
1386 assert_eq!(completed.ts_init, ts_init);
1387 }
1388
1389 #[rstest]
1390 fn test_stream_handler_completes_segmented_mcm_on_final_segment() {
1391 let ts_init = UnixNanos::from(1_800_000_000_000_000_005);
1392 let (handler, mut data_rx) = stream_handler_at(ts_init);
1393 let data = load_test_json("stream/mcm_SEGMENTS.jsonl");
1394 let mut segments = data.lines();
1395
1396 handler(stream_decode(segments.next().unwrap().as_bytes()).unwrap());
1397 handler(stream_decode(segments.next().unwrap().as_bytes()).unwrap());
1398
1399 assert!(data_rx.try_recv().is_err());
1400
1401 handler(stream_decode(segments.next().unwrap().as_bytes()).unwrap());
1402
1403 let custom = receive_custom::<BetfairSequenceCompleted>(&mut data_rx);
1404 let completed = custom
1405 .as_any()
1406 .downcast_ref::<BetfairSequenceCompleted>()
1407 .unwrap();
1408 assert_eq!(
1409 completed.ts_event,
1410 UnixNanos::from(1_700_000_000_000_000_000)
1411 );
1412 assert_eq!(completed.ts_init, ts_init);
1413 assert!(segments.next().is_none());
1414 assert!(data_rx.try_recv().is_err());
1415 }
1416
1417 #[rstest]
1418 fn test_stream_handler_stress_completes_each_segmented_mcm_once() {
1419 const SEQUENCE_COUNT: usize = 1_024;
1420 const MAX_MIDDLE_SEGMENTS: usize = 15;
1421
1422 let ts_init = UnixNanos::from(1_800_000_000_000_000_006);
1423 let (handler, mut data_rx) = stream_handler_at(ts_init);
1424 let data = load_test_json("stream/mcm_SEGMENTS.jsonl");
1425 let segments = data.lines().collect::<Vec<_>>();
1426
1427 for sequence in 0..SEQUENCE_COUNT {
1428 handler(stream_decode(segments[0].as_bytes()).unwrap());
1429 for _ in 0..sequence % (MAX_MIDDLE_SEGMENTS + 1) {
1430 handler(stream_decode(segments[1].as_bytes()).unwrap());
1431 }
1432
1433 assert!(data_rx.try_recv().is_err());
1434
1435 handler(stream_decode(segments[2].as_bytes()).unwrap());
1436
1437 let custom = receive_custom::<BetfairSequenceCompleted>(&mut data_rx);
1438 let completed = custom
1439 .as_any()
1440 .downcast_ref::<BetfairSequenceCompleted>()
1441 .unwrap();
1442 assert_eq!(
1443 completed.ts_event,
1444 UnixNanos::from(1_700_000_000_000_000_000)
1445 );
1446 assert_eq!(completed.ts_init, ts_init);
1447 assert!(data_rx.try_recv().is_err());
1448 }
1449 }
1450
1451 #[rstest]
1452 fn test_stream_handler_sets_rcm_init_from_clock() {
1453 let ts_init = UnixNanos::from(1_800_000_000_000_000_002);
1454
1455 let (handler, mut data_rx) = stream_handler_at(ts_init);
1456 let data = load_test_json("stream/rcm_single.json");
1457
1458 handler(stream_decode(data.as_bytes()).unwrap());
1459
1460 let custom = receive_custom::<BetfairRaceRunnerData>(&mut data_rx);
1461 let runner = custom
1462 .as_any()
1463 .downcast_ref::<BetfairRaceRunnerData>()
1464 .unwrap();
1465
1466 assert_eq!(runner.ts_event, UnixNanos::from(1_518_626_674_000_000_000));
1467 assert_eq!(runner.ts_init, ts_init);
1468 }
1469
1470 #[rstest]
1471 fn test_stream_handler_uses_rcm_publish_time_without_feed_time() {
1472 let ts_init = UnixNanos::from(1_800_000_000_000_000_003);
1473
1474 let (handler, mut data_rx) = stream_handler_at(ts_init);
1475 let data = load_test_json("stream/rcm_single.json");
1476 let mut message: serde_json::Value = serde_json::from_str(&data).unwrap();
1477 message
1478 .pointer_mut("/rc/0/rrc/0")
1479 .unwrap()
1480 .as_object_mut()
1481 .unwrap()
1482 .remove("ft");
1483 let data = message.to_string();
1484
1485 handler(stream_decode(data.as_bytes()).unwrap());
1486
1487 let custom = receive_custom::<BetfairRaceRunnerData>(&mut data_rx);
1488 let runner = custom
1489 .as_any()
1490 .downcast_ref::<BetfairRaceRunnerData>()
1491 .unwrap();
1492
1493 assert_eq!(runner.ts_event, UnixNanos::from(1_518_626_764_000_000_000));
1494 assert_eq!(runner.ts_init, ts_init);
1495 }
1496
1497 #[rstest]
1498 fn test_stream_handler_emits_cricket_match_custom_data() {
1499 let ts_init = UnixNanos::from(1_800_000_000_000_000_004);
1500
1501 let (handler, mut data_rx) = stream_handler_at(ts_init);
1502 let data = load_test_json("stream/ccm_single.json");
1503
1504 handler(stream_decode(data.as_bytes()).unwrap());
1505
1506 let event = data_rx.try_recv().expect("expected cricket custom data");
1507 let DataEvent::Data(Data::Custom(custom)) = event else {
1508 panic!("expected cricket custom data event, was {event:?}");
1509 };
1510 let cricket = custom
1511 .data
1512 .as_any()
1513 .downcast_ref::<BetfairCricketMatch>()
1514 .expect("custom data must be BetfairCricketMatch");
1515 let metadata = custom.data_type.metadata().expect("event metadata");
1516
1517 assert_eq!(cricket.event_id, "35741575");
1518 assert_eq!(cricket.market_id, "1.259334639");
1519 assert_eq!(cricket.ts_event, UnixNanos::from(1_700_000_000_000_000_000));
1520 assert_eq!(cricket.ts_init, ts_init);
1521 assert_eq!(
1522 metadata.get("event_id"),
1523 Some(&serde_json::Value::String("35741575".to_string())),
1524 );
1525 assert!(
1526 data_rx.try_recv().is_err(),
1527 "CCM fixture must emit exactly one event"
1528 );
1529 }
1530}