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