Skip to main content

nautilus_betfair/
data.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live market data client for the Betfair adapter.
17
18use 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
91/// Keep-alive interval in seconds (10 hours, matching Python default).
92const KEEP_ALIVE_INTERVAL_SECS: u64 = 36_000;
93
94/// Betfair live data client.
95#[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
121/// Wraps a custom data value with its instrument_id in both metadata (for
122/// topic routing) and identifier (for catalog partitioning).
123pub(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    /// Creates a new [`BetfairDataClient`] instance.
142    #[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        // Track cumulative traded volumes per (instrument_id, price) to compute
310        // incremental trade sizes. Betfair `trd` fields report totals, not deltas.
311        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                            // Emit instruments first so downstream consumers (DataEngine,
340                            // BacktestExchange) have the instrument cached before any status
341                            // or close event references it.
342                            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                        // Non-snapshot deltas and BSP deltas are buffered and flushed after
410                        // trades/tickers to mirror the Python `market_change_to_updates`
411                        // ordering (book deltas first, then BSP). Snapshots go inline.
412                        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        // Trades are included in market subscription via EX_TRADED
1168        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        // Instrument status is included in market subscription via EX_MARKET_DEF
1188        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        // Close transitions arrive via marketDefinition.status="CLOSED" on the
1208        // existing market subscription; no separate venue subscription exists.
1209        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}