Skip to main content

nautilus_betfair/
loader.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//! File-based loader for historical Betfair Exchange Streaming data.
17//!
18//! Reads compressed (gzip or bzip2) or plain JSON files containing
19//! newline-delimited Betfair ESA messages and produces Nautilus domain
20//! objects. The parsing logic mirrors the live data client handler in
21//! [`crate::data`].
22
23use std::{
24    fs::File,
25    io::{BufRead, BufReader},
26    path::Path,
27};
28
29use ahash::AHashMap;
30use anyhow::Context;
31use bzip2::read::BzDecoder;
32use flate2::read::GzDecoder;
33use nautilus_model::{
34    data::{InstrumentClose, InstrumentStatus, OrderBookDeltas, TradeTick},
35    identifiers::{InstrumentId, TradeId},
36    instruments::{Instrument, InstrumentAny},
37    types::{Currency, Money},
38};
39use rust_decimal::Decimal;
40
41use crate::{
42    common::{
43        enums::MarketStatus,
44        parse::{
45            make_instrument_id, parse_betfair_price, parse_betfair_quantity,
46            parse_market_definition, parse_millis_timestamp,
47        },
48    },
49    data_types::{
50        BetfairBspBookDelta, BetfairCricketMatch, BetfairRaceProgress, BetfairRaceRunnerData,
51        BetfairSequenceCompleted, BetfairStartingPrice, BetfairTicker,
52    },
53    stream::{
54        messages::{CCM, MCM, RCM, StreamMessage, stream_decode},
55        parse::{
56            make_trade_tick, parse_betfair_starting_prices, parse_betfair_ticker,
57            parse_bsp_book_deltas, parse_cricket_match, parse_instrument_closes,
58            parse_instrument_statuses, parse_race_progress, parse_race_runner_data,
59            parse_runner_book_deltas,
60        },
61    },
62};
63
64/// A parsed data item from a Betfair historical file.
65#[derive(Debug)]
66pub enum BetfairDataItem {
67    /// Instrument definition from a market definition.
68    Instrument(Box<InstrumentAny>),
69    /// Market status change for an instrument.
70    Status(InstrumentStatus),
71    /// Order book snapshot or delta update.
72    Deltas(OrderBookDeltas),
73    /// Incremental trade tick derived from cumulative traded volumes.
74    Trade(TradeTick),
75    /// Betfair-specific ticker data (last traded price, traded volume, BSP near/far).
76    Ticker(BetfairTicker),
77    /// Betfair Starting Price for a runner.
78    StartingPrice(BetfairStartingPrice),
79    /// BSP book delta (separate from exchange book).
80    BspBookDelta(BetfairBspBookDelta),
81    /// Instrument close event at market settlement.
82    InstrumentClose(InstrumentClose),
83    /// Marker emitted after each MCM batch is fully processed.
84    SequenceCompleted(BetfairSequenceCompleted),
85    /// GPS tracking data for a race runner (from RCM).
86    RaceRunnerData(BetfairRaceRunnerData),
87    /// Race-level progress data (from RCM).
88    RaceProgress(BetfairRaceProgress),
89    /// Cricket match data (from CCM).
90    CricketMatch(BetfairCricketMatch),
91}
92
93/// Reads Betfair historical data files and converts them into Nautilus domain objects.
94///
95/// Each file contains newline-delimited JSON from the Betfair Exchange Streaming API.
96/// The loader handles gzip decompression, stateful traded volume tracking, and
97/// instrument creation from market definitions.
98#[derive(Debug)]
99pub struct BetfairDataLoader {
100    currency: Currency,
101    min_notional: Option<Money>,
102    traded_volumes: AHashMap<(InstrumentId, Decimal), Decimal>,
103    instruments: AHashMap<InstrumentId, InstrumentAny>,
104}
105
106impl BetfairDataLoader {
107    /// Creates a new [`BetfairDataLoader`].
108    #[must_use]
109    pub fn new(currency: Currency, min_notional: Option<Money>) -> Self {
110        Self {
111            currency,
112            min_notional,
113            traded_volumes: AHashMap::new(),
114            instruments: AHashMap::new(),
115        }
116    }
117
118    /// Returns the instruments cached from the most recent load.
119    #[must_use]
120    pub fn instruments(&self) -> &AHashMap<InstrumentId, InstrumentAny> {
121        &self.instruments
122    }
123
124    /// Clears all cached state (instruments and traded volumes).
125    pub fn reset(&mut self) {
126        self.traded_volumes.clear();
127        self.instruments.clear();
128    }
129
130    /// Loads a Betfair historical data file and returns all parsed data items.
131    ///
132    /// Supports gzip-compressed (`.gz`), bzip2-compressed (`.bz2`), and plain JSON files.
133    /// Each line is deserialized as a Betfair stream message. MCM and RCM
134    /// messages are parsed into Nautilus domain objects; other message types
135    /// are skipped.
136    ///
137    /// # Errors
138    ///
139    /// Returns an error if the file cannot be opened or a line fails to parse.
140    pub fn load(&mut self, filepath: &Path) -> anyhow::Result<Vec<BetfairDataItem>> {
141        let reader = open_reader(filepath)?;
142        let mut items = Vec::new();
143
144        for (line_num, line_result) in reader.lines().enumerate() {
145            let line = line_result.with_context(|| {
146                format!(
147                    "failed to read line {} of '{}'",
148                    line_num + 1,
149                    filepath.display()
150                )
151            })?;
152
153            if line.is_empty() {
154                continue;
155            }
156
157            let msg = match stream_decode(line.as_bytes()) {
158                Ok(msg) => msg,
159                Err(e) => {
160                    log::warn!("Failed to decode line {}: {e}", line_num + 1);
161                    continue;
162                }
163            };
164
165            match msg {
166                StreamMessage::MarketChange(mcm) => self.process_mcm(&mcm, &mut items),
167                StreamMessage::RaceChange(rcm) => Self::process_rcm(&rcm, &mut items),
168                StreamMessage::CricketChange(ccm) => Self::process_ccm(&ccm, &mut items),
169                StreamMessage::Connection(_)
170                | StreamMessage::Status(_)
171                | StreamMessage::OrderChange(_) => {}
172            }
173        }
174
175        Ok(items)
176    }
177
178    /// Loads only instrument definitions from a Betfair historical data file.
179    ///
180    /// Scans the file for market definitions and creates instruments, but
181    /// skips all other data processing. Faster than `load()` when only
182    /// instruments are needed.
183    ///
184    /// # Errors
185    ///
186    /// Returns an error if the file cannot be opened or parsed.
187    pub fn load_instruments(&mut self, filepath: &Path) -> anyhow::Result<Vec<InstrumentAny>> {
188        let reader = open_reader(filepath)?;
189
190        for line_result in reader.lines() {
191            let line = line_result?;
192            if line.is_empty() {
193                continue;
194            }
195
196            let msg = match stream_decode(line.as_bytes()) {
197                Ok(msg) => msg,
198                Err(_) => continue,
199            };
200
201            if let StreamMessage::MarketChange(mcm) = msg {
202                let Some(market_changes) = &mcm.mc else {
203                    continue;
204                };
205
206                let ts_init = parse_millis_timestamp(mcm.pt);
207
208                for mc in market_changes {
209                    if let Some(def) = &mc.market_definition
210                        && let Ok(instruments) = parse_market_definition(
211                            &mc.id,
212                            def,
213                            self.currency,
214                            ts_init,
215                            ts_init,
216                            self.min_notional,
217                        )
218                    {
219                        for inst in instruments {
220                            self.instruments.insert(inst.id(), inst);
221                        }
222                    }
223                }
224            }
225        }
226
227        Ok(self.instruments.values().cloned().collect())
228    }
229
230    fn process_mcm(&mut self, mcm: &MCM, items: &mut Vec<BetfairDataItem>) {
231        if mcm.is_heartbeat() {
232            return;
233        }
234
235        let Some(market_changes) = &mcm.mc else {
236            return;
237        };
238
239        let ts_event = parse_millis_timestamp(mcm.pt);
240        let ts_init = ts_event;
241
242        for mc in market_changes {
243            let is_snapshot = mc.img;
244            let mut market_closed = false;
245
246            if let Some(def) = &mc.market_definition {
247                // Emit instruments first so sequential consumers (e.g. the backtest
248                // exchange) have the instrument in cache before any status or close
249                // event references it.
250                match parse_market_definition(
251                    &mc.id,
252                    def,
253                    self.currency,
254                    ts_event,
255                    ts_init,
256                    self.min_notional,
257                ) {
258                    Ok(new_instruments) => {
259                        for inst in &new_instruments {
260                            self.instruments.insert(inst.id(), inst.clone());
261                        }
262
263                        for inst in new_instruments {
264                            items.push(BetfairDataItem::Instrument(Box::new(inst)));
265                        }
266                    }
267                    Err(e) => {
268                        log::warn!("Failed to parse market definition for {}: {e}", mc.id);
269                    }
270                }
271
272                if let Some(status) = &def.status {
273                    market_closed = *status == MarketStatus::Closed;
274
275                    for event in parse_instrument_statuses(&mc.id, def, ts_event, ts_init) {
276                        items.push(BetfairDataItem::Status(event));
277                    }
278                }
279
280                for sp in parse_betfair_starting_prices(&mc.id, def, ts_event, ts_init) {
281                    items.push(BetfairDataItem::StartingPrice(sp));
282                }
283
284                for close in parse_instrument_closes(&mc.id, def, ts_event, ts_init) {
285                    items.push(BetfairDataItem::InstrumentClose(close));
286                }
287            }
288
289            // Non-snapshot deltas and BSP deltas are buffered and flushed after
290            // trades/tickers to mirror the Python `market_change_to_updates`
291            // ordering (book deltas first, then BSP). Snapshots go inline per
292            // runner, also matching Python.
293            let mut buffered_deltas: Vec<OrderBookDeltas> = Vec::new();
294            let mut buffered_bsp_deltas: Vec<BetfairBspBookDelta> = Vec::new();
295
296            if let Some(runner_changes) = &mc.rc {
297                for rc in runner_changes {
298                    let handicap = rc.hc.unwrap_or(Decimal::ZERO);
299                    let instrument_id = make_instrument_id(&mc.id, rc.id, handicap);
300
301                    match parse_runner_book_deltas(
302                        instrument_id,
303                        rc,
304                        is_snapshot,
305                        mcm.pt,
306                        ts_event,
307                        ts_init,
308                    ) {
309                        Ok(Some(deltas)) => {
310                            if is_snapshot {
311                                items.push(BetfairDataItem::Deltas(deltas));
312                            } else {
313                                buffered_deltas.push(deltas);
314                            }
315                        }
316                        Ok(None) => {}
317                        Err(e) => {
318                            log::warn!("Failed to parse book deltas for {instrument_id}: {e}");
319                        }
320                    }
321
322                    if let Some(trades) = &rc.trd {
323                        for pv in trades {
324                            if pv.volume == Decimal::ZERO {
325                                continue;
326                            }
327
328                            let key = (instrument_id, pv.price);
329                            let prev_volume = self
330                                .traded_volumes
331                                .get(&key)
332                                .copied()
333                                .unwrap_or(Decimal::ZERO);
334
335                            if pv.volume <= prev_volume {
336                                self.traded_volumes.insert(key, pv.volume);
337                                continue;
338                            }
339
340                            let trade_volume = pv.volume - prev_volume;
341                            self.traded_volumes.insert(key, pv.volume);
342
343                            let price = match parse_betfair_price(pv.price) {
344                                Ok(p) => p,
345                                Err(e) => {
346                                    log::warn!("Invalid trade price: {e}");
347                                    continue;
348                                }
349                            };
350                            let size = match parse_betfair_quantity(trade_volume) {
351                                Ok(q) => q,
352                                Err(e) => {
353                                    log::warn!("Invalid trade size: {e}");
354                                    continue;
355                                }
356                            };
357                            let trade_id =
358                                TradeId::new(format!("{}-{}-{}", mcm.pt, rc.id, pv.price));
359                            let tick = make_trade_tick(
360                                instrument_id,
361                                price,
362                                size,
363                                trade_id,
364                                ts_event,
365                                ts_init,
366                            );
367                            items.push(BetfairDataItem::Trade(tick));
368                        }
369                    }
370
371                    if let Some(ticker) = parse_betfair_ticker(instrument_id, rc, ts_event, ts_init)
372                    {
373                        items.push(BetfairDataItem::Ticker(ticker));
374                    }
375
376                    buffered_bsp_deltas.extend(parse_bsp_book_deltas(
377                        instrument_id,
378                        rc,
379                        ts_event,
380                        ts_init,
381                    ));
382                }
383            }
384
385            for deltas in buffered_deltas {
386                items.push(BetfairDataItem::Deltas(deltas));
387            }
388
389            for bsp_delta in buffered_bsp_deltas {
390                items.push(BetfairDataItem::BspBookDelta(bsp_delta));
391            }
392
393            if market_closed {
394                let prefix = format!("{}-", mc.id);
395                self.traded_volumes
396                    .retain(|k, _| !k.0.symbol.as_str().starts_with(&prefix));
397            }
398        }
399
400        items.push(BetfairDataItem::SequenceCompleted(
401            BetfairSequenceCompleted::new(ts_event, ts_init),
402        ));
403    }
404
405    fn process_rcm(rcm: &RCM, items: &mut Vec<BetfairDataItem>) {
406        let Some(race_changes) = &rcm.rc else {
407            return;
408        };
409
410        let ts_init = parse_millis_timestamp(rcm.pt);
411
412        for rc in race_changes {
413            let race_id = rc.id.as_deref().unwrap_or("");
414            let market_id = rc.mid.as_deref().unwrap_or("");
415
416            if let Some(runners) = &rc.rrc {
417                for rrc in runners {
418                    let ts_event = rrc.ft.map_or(ts_init, parse_millis_timestamp);
419
420                    if let Some(runner) =
421                        parse_race_runner_data(race_id, market_id, rrc, ts_event, ts_init)
422                    {
423                        items.push(BetfairDataItem::RaceRunnerData(runner));
424                    }
425                }
426            }
427
428            if let Some(rpc) = &rc.rpc {
429                let ts_event = rpc.ft.map_or(ts_init, parse_millis_timestamp);
430                let progress = parse_race_progress(race_id, market_id, rpc, ts_event, ts_init);
431                items.push(BetfairDataItem::RaceProgress(progress));
432            }
433        }
434    }
435
436    fn process_ccm(ccm: &CCM, items: &mut Vec<BetfairDataItem>) {
437        let Some(cricket_changes) = &ccm.cc else {
438            return;
439        };
440
441        let ts_init = parse_millis_timestamp(ccm.pt);
442
443        for cricket_change in cricket_changes {
444            if let Some(cricket) = parse_cricket_match(cricket_change, ts_init, ts_init) {
445                items.push(BetfairDataItem::CricketMatch(cricket));
446            }
447        }
448    }
449}
450
451fn open_reader(filepath: &Path) -> anyhow::Result<Box<dyn BufRead>> {
452    let file =
453        File::open(filepath).with_context(|| format!("failed to open '{}'", filepath.display()))?;
454
455    let ext = filepath.extension().and_then(|e| e.to_str()).unwrap_or("");
456
457    if ext.eq_ignore_ascii_case("gz") {
458        Ok(Box::new(BufReader::new(GzDecoder::new(file))))
459    } else if ext.eq_ignore_ascii_case("bz2") {
460        Ok(Box::new(BufReader::new(BzDecoder::new(file))))
461    } else {
462        Ok(Box::new(BufReader::new(file)))
463    }
464}
465
466#[cfg(test)]
467mod tests {
468    use std::path::PathBuf;
469
470    use nautilus_model::types::Price;
471    use rstest::rstest;
472
473    use super::*;
474    use crate::common::testing::load_test_json;
475
476    fn compact_json(pretty: &str) -> String {
477        let value: serde_json::Value = serde_json::from_str(pretty).unwrap();
478        serde_json::to_string(&value).unwrap()
479    }
480
481    fn local_data_dir() -> PathBuf {
482        PathBuf::from(env!("CARGO_MANIFEST_DIR"))
483            .ancestors()
484            .nth(3)
485            .unwrap()
486            .join("test_data/local/betfair")
487    }
488
489    fn test_data_dir() -> PathBuf {
490        PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("test_data")
491    }
492
493    #[rstest]
494    fn test_load_bz2_file() {
495        let filepath = test_data_dir().join("stream/sample.bz2");
496        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
497        let items = loader.load(&filepath).unwrap();
498
499        let instrument_count = items
500            .iter()
501            .filter(|i| matches!(i, BetfairDataItem::Instrument(_)))
502            .count();
503        assert!(
504            instrument_count > 0,
505            "should parse instruments from bz2 file"
506        );
507        assert_eq!(loader.instruments().len(), instrument_count);
508
509        let has_sequence = items
510            .iter()
511            .any(|i| matches!(i, BetfairDataItem::SequenceCompleted(_)));
512        assert!(has_sequence, "should emit SequenceCompleted");
513    }
514
515    #[rstest]
516    fn test_load_single_mcm_line() {
517        let data = compact_json(&load_test_json("stream/mcm_SUB_IMAGE.json"));
518        let tmp_dir = std::env::temp_dir().join("betfair_test");
519        std::fs::create_dir_all(&tmp_dir).unwrap();
520        let tmp_file = tmp_dir.join("test_single_mcm.json");
521        std::fs::write(&tmp_file, &data).unwrap();
522
523        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
524        let items = loader.load(&tmp_file).unwrap();
525
526        let instrument_count = items
527            .iter()
528            .filter(|i| matches!(i, BetfairDataItem::Instrument(_)))
529            .count();
530        assert!(
531            instrument_count > 0,
532            "should parse instruments from market definition"
533        );
534        assert_eq!(loader.instruments().len(), instrument_count);
535
536        let has_sequence = items
537            .iter()
538            .any(|i| matches!(i, BetfairDataItem::SequenceCompleted(_)));
539        assert!(has_sequence, "should emit SequenceCompleted");
540
541        std::fs::remove_file(&tmp_file).ok();
542    }
543
544    #[rstest]
545    fn test_load_rcm_uses_publish_time_as_ts_init() {
546        let rcm = r#"{"op":"rcm","clk":"17787241","pt":1583685064634,"rc":[{"id":"29741583.1630","mid":"1.169866199","rrc":[{"ft":1583685064600,"id":24330465,"long":-0.9116606,"lat":53.0712976,"spd":16.91,"prg":750.7,"sfq":2.52}],"rpc":{"ft":1583685064600,"g":"4f","st":11.74,"rt":30.93,"spd":16.8,"prg":749.2,"ord":[8709395,24330465]}}]}"#;
547        let tmp_dir = std::env::temp_dir().join("betfair_test");
548        std::fs::create_dir_all(&tmp_dir).unwrap();
549        let tmp_file = tmp_dir.join("test_rcm_timestamps.json");
550        std::fs::write(&tmp_file, rcm).unwrap();
551
552        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
553        let items = loader.load(&tmp_file).unwrap();
554
555        let expected_event = nautilus_core::UnixNanos::from(1_583_685_064_600_000_000);
556        let expected_init = nautilus_core::UnixNanos::from(1_583_685_064_634_000_000);
557
558        let runner = items
559            .iter()
560            .find_map(|item| match item {
561                BetfairDataItem::RaceRunnerData(runner) => Some(runner),
562                _ => None,
563            })
564            .expect("expected race runner data");
565        assert_eq!(runner.ts_event, expected_event);
566        assert_eq!(runner.ts_init, expected_init);
567
568        let progress = items
569            .iter()
570            .find_map(|item| match item {
571                BetfairDataItem::RaceProgress(progress) => Some(progress),
572                _ => None,
573            })
574            .expect("expected race progress data");
575        assert_eq!(progress.ts_event, expected_event);
576        assert_eq!(progress.ts_init, expected_init);
577
578        std::fs::remove_file(&tmp_file).ok();
579    }
580
581    #[rstest]
582    fn test_load_ccm_cricket_match() {
583        let ccm = compact_json(&load_test_json("stream/ccm_single.json"));
584        let tmp_dir = std::env::temp_dir().join("betfair_test");
585        std::fs::create_dir_all(&tmp_dir).unwrap();
586        let tmp_file = tmp_dir.join("test_ccm.json");
587        std::fs::write(&tmp_file, ccm).unwrap();
588
589        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
590        let items = loader.load(&tmp_file).unwrap();
591
592        let cricket = items
593            .iter()
594            .find_map(|item| match item {
595                BetfairDataItem::CricketMatch(cricket) => Some(cricket),
596                _ => None,
597            })
598            .expect("expected cricket match data");
599        assert_eq!(cricket.event_id, "35741575");
600        assert_eq!(cricket.market_id, "1.259334639");
601
602        std::fs::remove_file(&tmp_file).ok();
603    }
604
605    #[rstest]
606    fn test_load_mcm_with_book_data() {
607        let sub_image = compact_json(&load_test_json("stream/mcm_SUB_IMAGE.json"));
608        let update = compact_json(&load_test_json("stream/mcm_UPDATE.json"));
609
610        let tmp_dir = std::env::temp_dir().join("betfair_test");
611        std::fs::create_dir_all(&tmp_dir).unwrap();
612        let tmp_file = tmp_dir.join("test_book_data.json");
613        std::fs::write(&tmp_file, format!("{sub_image}\n{update}")).unwrap();
614
615        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
616        let items = loader.load(&tmp_file).unwrap();
617
618        let deltas_count = items
619            .iter()
620            .filter(|i| matches!(i, BetfairDataItem::Deltas(_)))
621            .count();
622        assert!(deltas_count > 0, "should parse book deltas");
623
624        std::fs::remove_file(&tmp_file).ok();
625    }
626
627    #[rstest]
628    fn test_load_instruments_only() {
629        let sub_image = compact_json(&load_test_json("stream/mcm_SUB_IMAGE.json"));
630        let update = compact_json(&load_test_json("stream/mcm_UPDATE.json"));
631
632        let tmp_dir = std::env::temp_dir().join("betfair_test");
633        std::fs::create_dir_all(&tmp_dir).unwrap();
634        let tmp_file = tmp_dir.join("test_instruments_only.json");
635        std::fs::write(&tmp_file, format!("{sub_image}\n{update}")).unwrap();
636
637        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
638        let instruments = loader.load_instruments(&tmp_file).unwrap();
639
640        assert!(!instruments.is_empty(), "should find instruments");
641        assert_eq!(loader.instruments().len(), instruments.len());
642
643        std::fs::remove_file(&tmp_file).ok();
644    }
645
646    #[rstest]
647    fn test_reset_clears_state() {
648        let data = compact_json(&load_test_json("stream/mcm_SUB_IMAGE.json"));
649        let tmp_dir = std::env::temp_dir().join("betfair_test");
650        std::fs::create_dir_all(&tmp_dir).unwrap();
651        let tmp_file = tmp_dir.join("test_reset.json");
652        std::fs::write(&tmp_file, &data).unwrap();
653
654        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
655        loader.load(&tmp_file).unwrap();
656        assert!(!loader.instruments().is_empty());
657
658        loader.reset();
659        assert!(loader.instruments().is_empty());
660        assert!(loader.traded_volumes.is_empty());
661
662        std::fs::remove_file(&tmp_file).ok();
663    }
664
665    #[rstest]
666    fn test_load_bsp_data() {
667        let raw = load_test_json("stream/mcm_BSP.json");
668        let messages: Vec<serde_json::Value> = serde_json::from_str(&raw).unwrap();
669        let lines: Vec<String> = messages
670            .iter()
671            .map(|v| serde_json::to_string(v).unwrap())
672            .collect();
673
674        let tmp_dir = std::env::temp_dir().join("betfair_test");
675        std::fs::create_dir_all(&tmp_dir).unwrap();
676        let tmp_file = tmp_dir.join("test_bsp.json");
677        std::fs::write(&tmp_file, lines.join("\n")).unwrap();
678
679        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
680        let items = loader.load(&tmp_file).unwrap();
681
682        let bsp_count = items
683            .iter()
684            .filter(|i| matches!(i, BetfairDataItem::BspBookDelta(_)))
685            .count();
686        assert!(bsp_count > 0, "should parse BSP book deltas");
687
688        std::fs::remove_file(&tmp_file).ok();
689    }
690
691    #[rstest]
692    fn test_load_market_definition_with_traded_volumes() {
693        let sub_image = compact_json(&load_test_json("stream/mcm_SUB_IMAGE.json"));
694        let update_tv = compact_json(&load_test_json("stream/mcm_UPDATE_tv.json"));
695
696        let tmp_dir = std::env::temp_dir().join("betfair_test");
697        std::fs::create_dir_all(&tmp_dir).unwrap();
698        let tmp_file = tmp_dir.join("test_tv.json");
699        std::fs::write(&tmp_file, format!("{sub_image}\n{update_tv}")).unwrap();
700
701        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
702        let items = loader.load(&tmp_file).unwrap();
703
704        let ticker_count = items
705            .iter()
706            .filter(|i| matches!(i, BetfairDataItem::Ticker(_)))
707            .count();
708        assert!(ticker_count > 0, "should parse ticker data from tv updates");
709
710        std::fs::remove_file(&tmp_file).ok();
711    }
712
713    #[rstest]
714    #[ignore] // Requires user-fetched data in test_data/local/betfair/
715    fn test_load_match_odds_file() {
716        let filepath = local_data_dir().join("1.253378068.gz");
717        if !filepath.exists() {
718            eprintln!("Skipping: {filepath:?} not found");
719            return;
720        }
721
722        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
723        let items = loader.load(&filepath).unwrap();
724
725        let instrument_count = items
726            .iter()
727            .filter(|i| matches!(i, BetfairDataItem::Instrument(_)))
728            .count();
729        let deltas_count = items
730            .iter()
731            .filter(|i| matches!(i, BetfairDataItem::Deltas(_)))
732            .count();
733        let trade_count = items
734            .iter()
735            .filter(|i| matches!(i, BetfairDataItem::Trade(_)))
736            .count();
737        let close_count = items
738            .iter()
739            .filter(|i| matches!(i, BetfairDataItem::InstrumentClose(_)))
740            .count();
741
742        println!(
743            "Match odds file: {instrument_count} instruments, {deltas_count} deltas, {trade_count} trades, {close_count} closes"
744        );
745        println!("Total items: {}", items.len());
746
747        // 3 runners (home/draw/away), emitted on each market definition
748        assert!(instrument_count >= 3, "expected at least 3 instruments");
749        assert!(deltas_count > 0, "expected book deltas");
750        assert!(trade_count > 0, "expected trade ticks");
751        assert!(close_count > 0, "expected instrument closes at settlement");
752
753        // Winner should be runner 2426
754        let closes: Vec<_> = items
755            .iter()
756            .filter_map(|i| match i {
757                BetfairDataItem::InstrumentClose(c) => Some(c),
758                _ => None,
759            })
760            .collect();
761        let winner = closes.iter().find(|c| c.close_price == Price::from("1.00"));
762        assert!(winner.is_some(), "expected a winner with close_price 1.00");
763        assert!(
764            winner
765                .unwrap()
766                .instrument_id
767                .symbol
768                .as_str()
769                .contains("2426"),
770            "winner should be runner 2426"
771        );
772    }
773
774    #[rstest]
775    #[ignore] // Requires user-fetched data in test_data/local/betfair/
776    fn test_load_racing_win_file() {
777        let filepath = local_data_dir().join("1.245077076.gz");
778        if !filepath.exists() {
779            eprintln!("Skipping: {filepath:?} not found");
780            return;
781        }
782
783        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
784        let items = loader.load(&filepath).unwrap();
785
786        let instrument_count = items
787            .iter()
788            .filter(|i| matches!(i, BetfairDataItem::Instrument(_)))
789            .count();
790        let deltas_count = items
791            .iter()
792            .filter(|i| matches!(i, BetfairDataItem::Deltas(_)))
793            .count();
794        let trade_count = items
795            .iter()
796            .filter(|i| matches!(i, BetfairDataItem::Trade(_)))
797            .count();
798        let close_count = items
799            .iter()
800            .filter(|i| matches!(i, BetfairDataItem::InstrumentClose(_)))
801            .count();
802
803        println!(
804            "Racing file: {instrument_count} instruments, {deltas_count} deltas, {trade_count} trades, {close_count} closes"
805        );
806        println!("Total items: {}", items.len());
807
808        // 6 runners (though 2 removed during the race)
809        assert!(instrument_count >= 6, "expected at least 6 instruments");
810        assert!(deltas_count > 0, "expected book deltas");
811        assert!(trade_count > 0, "expected trade ticks");
812        assert!(close_count > 0, "expected instrument closes at settlement");
813
814        // Winner should be runner 75925986
815        let closes: Vec<_> = items
816            .iter()
817            .filter_map(|i| match i {
818                BetfairDataItem::InstrumentClose(c) => Some(c),
819                _ => None,
820            })
821            .collect();
822        let winner = closes.iter().find(|c| c.close_price == Price::from("1.00"));
823        assert!(winner.is_some(), "expected a winner with close_price 1.00");
824        assert!(
825            winner
826                .unwrap()
827                .instrument_id
828                .symbol
829                .as_str()
830                .contains("75925986"),
831            "winner should be runner 75925986"
832        );
833    }
834
835    fn write_tmp(contents: &str, name: &str) -> PathBuf {
836        let tmp_dir = std::env::temp_dir().join("betfair_test");
837        std::fs::create_dir_all(&tmp_dir).unwrap();
838        let tmp_file = tmp_dir.join(name);
839        std::fs::write(&tmp_file, contents).unwrap();
840        tmp_file
841    }
842
843    fn find_first(
844        items: &[BetfairDataItem],
845        pred: impl Fn(&BetfairDataItem) -> bool,
846    ) -> Option<usize> {
847        items.iter().position(pred)
848    }
849
850    fn find_last(
851        items: &[BetfairDataItem],
852        pred: impl Fn(&BetfairDataItem) -> bool,
853    ) -> Option<usize> {
854        items.iter().rposition(pred)
855    }
856
857    /// Split the loader output into per-MCM slices using `SequenceCompleted`
858    /// as the delimiter. Each MCM ends with one `SequenceCompleted` item.
859    fn partition_by_mcm(items: &[BetfairDataItem]) -> Vec<&[BetfairDataItem]> {
860        let mut partitions = Vec::new();
861        let mut start = 0;
862
863        for (i, item) in items.iter().enumerate() {
864            if matches!(item, BetfairDataItem::SequenceCompleted(_)) {
865                partitions.push(&items[start..=i]);
866                start = i + 1;
867            }
868        }
869
870        partitions
871    }
872
873    #[rstest]
874    fn test_load_emits_instrument_before_status_and_close() {
875        // The loader must emit `Instrument` events before any `InstrumentStatus`
876        // or `InstrumentClose` in the same MCM so downstream consumers (e.g. the
877        // backtest exchange) have the instrument cached before lifecycle events
878        // are processed.
879        let data = compact_json(&load_test_json("stream/mcm_UPDATE_md.json"));
880        let tmp_file = write_tmp(&data, "test_order_instrument_first.json");
881
882        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
883        let items = loader.load(&tmp_file).unwrap();
884
885        let instrument_idx =
886            find_first(&items, |i| matches!(i, BetfairDataItem::Instrument(_))).unwrap();
887        let status_idx = find_first(&items, |i| matches!(i, BetfairDataItem::Status(_))).unwrap();
888
889        assert!(
890            instrument_idx < status_idx,
891            "Instrument (idx {instrument_idx}) must precede Status (idx {status_idx})"
892        );
893
894        std::fs::remove_file(&tmp_file).ok();
895    }
896
897    #[rstest]
898    fn test_load_emits_instrument_before_close() {
899        // Synthetic fixture: market CLOSED with terminal runner statuses so that
900        // both Instrument and InstrumentClose are emitted within the same MCM.
901        // Instrument must appear first.
902        let mcm = r#"{"op":"mcm","id":1,"pt":1627617202953,"ct":"SUB_IMAGE","mc":[{"id":"1.1","marketDefinition":{"bspMarket":false,"turnInPlayEnabled":false,"persistenceEnabled":false,"marketBaseRate":5,"eventId":"1","eventTypeId":"1","numberOfWinners":1,"bettingType":"ODDS","marketType":"WIN","marketTime":"2021-07-30T03:55:00.000Z","bspReconciled":true,"complete":true,"inPlay":false,"crossMatching":false,"runnersVoidable":false,"numberOfActiveRunners":0,"betDelay":0,"status":"CLOSED","runners":[{"status":"WINNER","sortPriority":1,"id":101},{"status":"LOSER","sortPriority":2,"id":102}],"regulators":["MR_INT"],"discountAllowed":true,"timezone":"UTC","openDate":"2021-07-30T02:45:00.000Z","version":1,"priceLadderDefinition":{"type":"CLASSIC"}}}]}"#;
903        let tmp_file = write_tmp(mcm, "test_order_instrument_before_close.json");
904
905        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
906        let items = loader.load(&tmp_file).unwrap();
907
908        let instrument_idx =
909            find_first(&items, |i| matches!(i, BetfairDataItem::Instrument(_))).unwrap();
910        let close_idx =
911            find_first(&items, |i| matches!(i, BetfairDataItem::InstrumentClose(_))).unwrap();
912
913        assert!(
914            instrument_idx < close_idx,
915            "Instrument (idx {instrument_idx}) must precede InstrumentClose (idx {close_idx})"
916        );
917
918        std::fs::remove_file(&tmp_file).ok();
919    }
920
921    #[rstest]
922    fn test_load_non_snapshot_deltas_tail_after_trades() {
923        // Non-snapshot runner updates must emit book deltas AFTER any trades or
924        // tickers parsed from the same message, matching the Python
925        // `market_change_to_updates` ordering and keeping live/backtest in step.
926        let data = compact_json(&load_test_json("stream/mcm_live_UPDATE.json"));
927        let tmp_file = write_tmp(&data, "test_order_deltas_tail.json");
928
929        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
930        let items = loader.load(&tmp_file).unwrap();
931
932        // Fixture is a single non-snapshot MCM (img=None) with trd + atl on one rc
933        let last_trade_idx = find_last(&items, |i| matches!(i, BetfairDataItem::Trade(_))).unwrap();
934        let first_deltas_idx =
935            find_first(&items, |i| matches!(i, BetfairDataItem::Deltas(_))).unwrap();
936
937        assert!(
938            first_deltas_idx > last_trade_idx,
939            "Deltas (first idx {first_deltas_idx}) must tail after Trade (last idx {last_trade_idx}) on non-snapshot updates"
940        );
941
942        std::fs::remove_file(&tmp_file).ok();
943    }
944
945    #[rstest]
946    fn test_load_snapshot_deltas_emit_inline_before_trades() {
947        // Snapshot messages (mc.img=true) emit Clear+Add deltas inline per runner
948        // so consumers can apply the book state before any trades in the same
949        // MCM. This matches Python's inline-snapshot behaviour.
950        let data = compact_json(&load_test_json(
951            "stream/market_definition_runner_removed.json",
952        ));
953        let tmp_file = write_tmp(&data, "test_order_snapshot_inline.json");
954
955        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
956        let items = loader.load(&tmp_file).unwrap();
957
958        let first_deltas_idx =
959            find_first(&items, |i| matches!(i, BetfairDataItem::Deltas(_))).unwrap();
960        let first_trade_idx =
961            find_first(&items, |i| matches!(i, BetfairDataItem::Trade(_))).unwrap();
962
963        assert!(
964            first_deltas_idx < first_trade_idx,
965            "Snapshot Deltas (first idx {first_deltas_idx}) must emit before Trade (first idx {first_trade_idx})"
966        );
967
968        std::fs::remove_file(&tmp_file).ok();
969    }
970
971    #[rstest]
972    fn test_load_bsp_tails_after_book_deltas() {
973        // Within each MCM, BSP deltas must emit after all regular book deltas.
974        // Python flushes `book_updates` before `bsp_book_updates`; the Rust
975        // loader must do the same to preserve consumer ordering.
976        let raw = load_test_json("stream/mcm_BSP.json");
977        let messages: Vec<serde_json::Value> = serde_json::from_str(&raw).unwrap();
978        let lines: Vec<String> = messages
979            .iter()
980            .map(|v| serde_json::to_string(v).unwrap())
981            .collect();
982        let tmp_file = write_tmp(&lines.join("\n"), "test_order_bsp_tail.json");
983
984        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
985        let items = loader.load(&tmp_file).unwrap();
986
987        let partitions = partition_by_mcm(&items);
988        assert!(
989            !partitions.is_empty(),
990            "expected at least one MCM partition"
991        );
992
993        let mut checked_any = false;
994
995        for partition in partitions {
996            let last_deltas_idx = find_last(partition, |i| matches!(i, BetfairDataItem::Deltas(_)));
997            let first_bsp_idx =
998                find_first(partition, |i| matches!(i, BetfairDataItem::BspBookDelta(_)));
999
1000            if let (Some(last_deltas), Some(first_bsp)) = (last_deltas_idx, first_bsp_idx) {
1001                assert!(
1002                    last_deltas < first_bsp,
1003                    "BspBookDelta (first idx {first_bsp}) must tail after Deltas (last idx {last_deltas}) within the same MCM"
1004                );
1005                checked_any = true;
1006            }
1007        }
1008
1009        assert!(
1010            checked_any,
1011            "expected at least one MCM to contain both Deltas and BspBookDelta"
1012        );
1013
1014        std::fs::remove_file(&tmp_file).ok();
1015    }
1016
1017    #[rstest]
1018    fn test_load_emits_close_for_removed_runner_while_market_open() {
1019        // Removed runners must fire InstrumentClose as soon as the market
1020        // definition reports the Removed status, regardless of whether the
1021        // market as a whole is still Open. This matches the Python parser.
1022        let mcm = r#"{"op":"mcm","id":1,"pt":1627617202953,"ct":"SUB_IMAGE","mc":[{"id":"1.2","marketDefinition":{"bspMarket":false,"turnInPlayEnabled":false,"persistenceEnabled":false,"marketBaseRate":5,"eventId":"1","eventTypeId":"1","numberOfWinners":1,"bettingType":"ODDS","marketType":"WIN","marketTime":"2021-07-30T03:55:00.000Z","bspReconciled":false,"complete":true,"inPlay":false,"crossMatching":false,"runnersVoidable":false,"numberOfActiveRunners":1,"betDelay":0,"status":"OPEN","runners":[{"status":"ACTIVE","sortPriority":1,"id":201},{"status":"REMOVED","sortPriority":2,"id":202}],"regulators":["MR_INT"],"discountAllowed":true,"timezone":"UTC","openDate":"2021-07-30T02:45:00.000Z","version":1,"priceLadderDefinition":{"type":"CLASSIC"}}}]}"#;
1023        let tmp_file = write_tmp(mcm, "test_order_close_for_removed.json");
1024
1025        let mut loader = BetfairDataLoader::new(Currency::GBP(), None);
1026        let items = loader.load(&tmp_file).unwrap();
1027
1028        let closes: Vec<_> = items
1029            .iter()
1030            .filter_map(|i| match i {
1031                BetfairDataItem::InstrumentClose(c) => Some(c),
1032                _ => None,
1033            })
1034            .collect();
1035
1036        assert_eq!(
1037            closes.len(),
1038            1,
1039            "Removed runner must produce exactly one InstrumentClose while market is Open"
1040        );
1041        assert!(
1042            closes[0].instrument_id.symbol.as_str().contains("202"),
1043            "close must target the removed runner (selection id 202)"
1044        );
1045
1046        std::fs::remove_file(&tmp_file).ok();
1047    }
1048}