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