Skip to main content

nautilus_betfair/stream/
parse.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//! Parsing utilities that convert Betfair stream messages into Nautilus domain models.
17//!
18//! MCM (Market Change Messages) are parsed into order book deltas, trade ticks,
19//! and instrument status updates. OCM (Order Change Messages) are parsed into
20//! order status reports and fill reports.
21
22use ahash::{AHashMap, AHashSet};
23use nautilus_core::UnixNanos;
24use nautilus_model::{
25    data::{
26        BookOrder, InstrumentClose, InstrumentStatus, OrderBookDelta, OrderBookDeltas, TradeTick,
27    },
28    enums::{
29        AggressorSide, BookAction, InstrumentCloseType, LiquiditySide, MarketStatusAction,
30        OrderSide, OrderStatus, OrderType, RecordFlag, TimeInForce,
31    },
32    identifiers::{AccountId, ClientOrderId, InstrumentId, TradeId, VenueOrderId},
33    reports::{FillReport, OrderStatusReport},
34    types::{Currency, Money, Price, Quantity},
35};
36use rust_decimal::Decimal;
37
38use crate::{
39    common::{
40        enums::{
41            MarketStatus, RunnerStatus, StreamingOrderStatus, StreamingOrderType,
42            resolve_streaming_order_status,
43        },
44        parse::{
45            make_instrument_id, normalize_betfair_price, normalize_betfair_quantity,
46            parse_betfair_price, parse_betfair_quantity, parse_millis_timestamp,
47        },
48    },
49    data_types::{
50        BetfairBspBookDelta, BetfairCricketMatch, BetfairRaceProgress, BetfairRaceRunnerData,
51        BetfairStartingPrice, BetfairTicker,
52    },
53    stream::messages::{
54        CricketChange, MarketDefinition, RaceProgressChange, RaceRunnerChange, RunnerChange,
55        UnmatchedOrder,
56    },
57};
58
59/// Parses a single runner's book data into [`OrderBookDeltas`].
60///
61/// Handles both full image snapshots (`is_snapshot = true`) and delta updates.
62///
63/// Only processes price-keyed fields (`atb`/`atl`). Level-indexed fields
64/// (`batb`/`batl`/`bdatb`/`bdatl`) are ignored because correct delta
65/// processing requires stateful level-to-price tracking. The stream
66/// subscription should use `EX_ALL_OFFERS` which populates `atb`/`atl`.
67///
68/// Book side mapping (Betfair exchange convention):
69/// - `atb` (available to back) -> [`OrderSide::Buy`] (bid side)
70/// - `atl` (available to lay) -> [`OrderSide::Sell`] (ask side)
71///
72/// For snapshots: emits a Clear delta followed by Add deltas for each level.
73/// For updates: emits Update or Delete (when volume is zero) deltas.
74///
75/// Returns `Ok(None)` if the runner change contains no processable book data.
76///
77/// # Errors
78///
79/// Returns an error if price or quantity values cannot be converted.
80pub fn parse_runner_book_deltas(
81    instrument_id: InstrumentId,
82    rc: &RunnerChange,
83    is_snapshot: bool,
84    sequence: u64,
85    ts_event: UnixNanos,
86    ts_init: UnixNanos,
87) -> anyhow::Result<Option<OrderBookDeltas>> {
88    let atb_len = rc.atb.as_ref().map_or(0, Vec::len);
89    let atl_len = rc.atl.as_ref().map_or(0, Vec::len);
90    let total_levels = atb_len + atl_len;
91
92    if total_levels == 0 && !is_snapshot {
93        return Ok(None);
94    }
95
96    let snapshot_flags = if is_snapshot {
97        RecordFlag::F_SNAPSHOT as u8
98    } else {
99        0
100    };
101    let mut deltas = Vec::with_capacity(total_levels + usize::from(is_snapshot));
102
103    if is_snapshot {
104        let mut clear = OrderBookDelta::clear(instrument_id, sequence, ts_event, ts_init);
105
106        if total_levels == 0 {
107            clear.flags |= RecordFlag::F_LAST as u8;
108        }
109        deltas.push(clear);
110    }
111
112    for (levels, side) in [
113        (rc.atb.as_deref().unwrap_or(&[]), OrderSide::Buy),
114        (rc.atl.as_deref().unwrap_or(&[]), OrderSide::Sell),
115    ] {
116        for pv in levels {
117            if is_snapshot && pv.volume == Decimal::ZERO {
118                continue;
119            }
120
121            let action = if is_snapshot {
122                BookAction::Add
123            } else if pv.volume == Decimal::ZERO {
124                BookAction::Delete
125            } else {
126                BookAction::Update
127            };
128
129            deltas.push(OrderBookDelta::new(
130                instrument_id,
131                action,
132                BookOrder::new(
133                    side,
134                    parse_betfair_price(pv.price)?,
135                    parse_betfair_quantity(pv.volume)?,
136                    0,
137                ),
138                snapshot_flags,
139                sequence,
140                ts_event,
141                ts_init,
142            ));
143        }
144    }
145
146    // Set F_LAST on the final delta
147    if let Some(last) = deltas.last_mut() {
148        last.flags |= RecordFlag::F_LAST as u8;
149    }
150
151    Ok(Some(OrderBookDeltas::new(instrument_id, deltas)))
152}
153
154/// Creates a [`TradeTick`] from stream data.
155///
156/// Betfair does not identify the aggressor side in its stream, so
157/// [`AggressorSide::NoAggressor`] is always used.
158#[must_use]
159pub fn make_trade_tick(
160    instrument_id: InstrumentId,
161    price: Price,
162    size: Quantity,
163    trade_id: TradeId,
164    ts_event: UnixNanos,
165    ts_init: UnixNanos,
166) -> TradeTick {
167    TradeTick::new(
168        instrument_id,
169        price,
170        size,
171        AggressorSide::NoAggressor,
172        trade_id,
173        ts_event,
174        ts_init,
175    )
176}
177
178/// Produces per-runner [`InstrumentStatus`] events from a market definition.
179///
180/// Iterates `def.runners` and maps each runner's lifecycle to a Nautilus status.
181/// Scratched runners (`Removed`, `RemovedVacant`) close immediately regardless
182/// of market-level state. The `in_play` flag distinguishes pre-open (Open + not
183/// in play) from active trading (Open + in play).
184///
185/// Returns an empty vector when `def.status` or `def.runners` is missing.
186#[must_use]
187pub fn parse_instrument_statuses(
188    market_id: &str,
189    def: &MarketDefinition,
190    ts_event: UnixNanos,
191    ts_init: UnixNanos,
192) -> Vec<InstrumentStatus> {
193    let Some(status) = def.status else {
194        return Vec::new();
195    };
196    let Some(runners) = &def.runners else {
197        return Vec::new();
198    };
199    let in_play = def.in_play.unwrap_or(false);
200
201    if status == MarketStatus::Unknown {
202        log::warn!("Skipping unmodeled Betfair market status for market {market_id}");
203        return Vec::new();
204    }
205
206    runners
207        .iter()
208        .filter_map(|rd| {
209            let handicap = rd.hc.unwrap_or(Decimal::ZERO);
210            let instrument_id = make_instrument_id(market_id, rd.id, handicap);
211            if rd.status == Some(RunnerStatus::Unknown) {
212                log::warn!("Skipping unmodeled Betfair runner status for {instrument_id}");
213                return None;
214            }
215            let action = match rd.status {
216                Some(RunnerStatus::Removed | RunnerStatus::RemovedVacant) => {
217                    MarketStatusAction::Close
218                }
219                _ => match (status, in_play) {
220                    (MarketStatus::Inactive, _) => MarketStatusAction::Close,
221                    (MarketStatus::Open, false) => MarketStatusAction::PreOpen,
222                    (MarketStatus::Open, true) => MarketStatusAction::Trading,
223                    (MarketStatus::Suspended, _) => MarketStatusAction::Pause,
224                    (MarketStatus::Closed, _) => MarketStatusAction::Close,
225                    // Unreachable: unmodeled market status skips the whole definition above.
226                    (MarketStatus::Unknown, _) => MarketStatusAction::None,
227                },
228            };
229            let is_trading = action == MarketStatusAction::Trading;
230            Some(InstrumentStatus::new(
231                instrument_id,
232                action,
233                ts_event,
234                ts_init,
235                None,
236                None,
237                Some(is_trading),
238                None,
239                None,
240            ))
241        })
242        .collect()
243}
244
245/// Generates a deterministic [`TradeId`] for a Betfair fill.
246///
247/// Uses `bet_id` and cumulative `sm` (size matched) which together uniquely
248/// identify each fill state, since `sm` increases monotonically with each fill.
249pub fn make_trade_id(uo: &UnmatchedOrder) -> TradeId {
250    let sm = normalize_betfair_quantity(uo.sm.unwrap_or(Decimal::ZERO));
251    make_trade_id_for_size(&uo.id, sm)
252}
253
254/// Tracks cumulative fill state per bet to compute incremental fills from the
255/// Betfair OCM stream.
256///
257/// Betfair provides cumulative `sm` (size matched) and `avp` (average price
258/// matched) on each order update. This tracker maintains per-bet state to
259/// derive individual fill quantities and prices for each update.
260#[derive(Clone, Debug, Default)]
261pub struct FillTracker {
262    filled_qty: AHashMap<String, Decimal>,
263    voided_qty: AHashMap<String, Decimal>,
264    avg_px: AHashMap<String, Decimal>,
265    published_trade_ids: AHashSet<String>,
266    fill_lots: AHashMap<String, Vec<FillLot>>,
267    fill_voids: AHashMap<(String, TradeId), Decimal>,
268}
269
270#[derive(Debug, Clone)]
271struct FillLot {
272    trade_id: TradeId,
273    quantity: Decimal,
274    price: Price,
275}
276
277/// One cumulative per-fill allocation derived from Betfair's cumulative `sv`.
278#[derive(Debug, Clone)]
279pub struct FillVoidAllocation {
280    pub trade_id: TradeId,
281    pub voided_qty: Quantity,
282    pub last_px: Price,
283}
284
285impl FillTracker {
286    /// Creates a new [`FillTracker`] instance.
287    #[must_use]
288    pub fn new() -> Self {
289        Self::default()
290    }
291
292    /// Computes an incremental [`FillReport`] for an unmatched order update.
293    ///
294    /// Returns `None` if no new fill occurred (size matched unchanged,
295    /// duplicate trade ID, or overfill detected).
296    #[expect(clippy::too_many_arguments)]
297    pub fn maybe_fill_report(
298        &mut self,
299        uo: &UnmatchedOrder,
300        order_qty: Decimal,
301        instrument_id: InstrumentId,
302        account_id: AccountId,
303        currency: Currency,
304        ts_event: UnixNanos,
305        ts_init: UnixNanos,
306    ) -> Option<FillReport> {
307        let raw_sm = uo.sm?;
308        let sm = normalize_betfair_quantity(raw_sm);
309        let order_qty = normalize_betfair_quantity(resolve_stream_order_quantity(order_qty, uo));
310
311        if sm > order_qty {
312            log::warn!(
313                "Rejecting potential overfill for bet_id={}: order_qty={order_qty}, sm={sm}",
314                uo.id,
315            );
316            return None;
317        }
318
319        let (trade_id, last_qty, last_px) =
320            if uo.sv.is_some_and(|sv| sv > Decimal::ZERO) && self.has_fill_lots(&uo.id) {
321                self.advance_cumulative_fill_with_voids(
322                    &uo.id,
323                    raw_sm,
324                    uo.sv.unwrap_or(Decimal::ZERO),
325                    uo.avp,
326                    uo.p,
327                )?
328            } else {
329                self.advance_cumulative_fill(&uo.id, raw_sm, uo.avp, uo.p)?
330            };
331
332        let venue_order_id = VenueOrderId::from(uo.id.as_str());
333        let order_side = OrderSide::from(uo.side);
334        let client_order_id = uo
335            .rfo
336            .as_deref()
337            .filter(|s| !s.is_empty())
338            .map(ClientOrderId::from);
339        let ts_fill = uo.md.map_or(ts_event, parse_millis_timestamp);
340
341        Some(make_fill_report(
342            account_id,
343            instrument_id,
344            venue_order_id,
345            trade_id,
346            order_side,
347            last_qty,
348            last_px,
349            currency,
350            client_order_id,
351            ts_fill,
352            ts_init,
353        ))
354    }
355
356    pub(crate) fn advance_cumulative_fill(
357        &mut self,
358        bet_id: &str,
359        raw_size_matched: Decimal,
360        average_price_matched: Option<Decimal>,
361        fallback_price: Decimal,
362    ) -> Option<(TradeId, Quantity, Price)> {
363        self.advance_cumulative_fill_inner(
364            bet_id,
365            raw_size_matched,
366            None,
367            average_price_matched,
368            fallback_price,
369        )
370    }
371
372    pub(crate) fn advance_cumulative_fill_with_voids(
373        &mut self,
374        bet_id: &str,
375        raw_size_matched: Decimal,
376        cumulative_voided: Decimal,
377        average_price_matched: Option<Decimal>,
378        fallback_price: Decimal,
379    ) -> Option<(TradeId, Quantity, Price)> {
380        self.advance_cumulative_fill_inner(
381            bet_id,
382            raw_size_matched,
383            Some(cumulative_voided),
384            average_price_matched,
385            fallback_price,
386        )
387    }
388
389    fn advance_cumulative_fill_inner(
390        &mut self,
391        bet_id: &str,
392        raw_size_matched: Decimal,
393        cumulative_voided: Option<Decimal>,
394        average_price_matched: Option<Decimal>,
395        fallback_price: Decimal,
396    ) -> Option<(TradeId, Quantity, Price)> {
397        let size_matched = normalize_betfair_quantity(raw_size_matched);
398        if size_matched <= Decimal::ZERO {
399            return None;
400        }
401
402        let previous_filled = self
403            .filled_qty
404            .get(bet_id)
405            .copied()
406            .map_or(Decimal::ZERO, normalize_betfair_quantity);
407
408        if size_matched == previous_filled {
409            self.record_cumulative_state(bet_id, size_matched, average_price_matched);
410            return None;
411        }
412
413        if size_matched < previous_filled {
414            return None;
415        }
416
417        let trade_id = make_trade_id_for_size(bet_id, size_matched);
418        let raw_trade_id = make_trade_id_for_size(bet_id, raw_size_matched);
419        if self.published_trade_ids.contains(trade_id.as_str())
420            || self.published_trade_ids.contains(raw_trade_id.as_str())
421        {
422            self.record_cumulative_state(bet_id, size_matched, average_price_matched);
423            self.published_trade_ids.insert(trade_id.to_string());
424            return None;
425        }
426
427        let fill_qty = size_matched - previous_filled;
428        let fill_price = if let Some(cumulative_voided) = cumulative_voided {
429            self.compute_fill_price_with_voids(
430                bet_id,
431                average_price_matched,
432                fallback_price,
433                size_matched,
434                fill_qty,
435                cumulative_voided,
436            )
437        } else {
438            self.compute_fill_price(
439                bet_id,
440                average_price_matched,
441                fallback_price,
442                previous_filled,
443                size_matched,
444            )
445        };
446        let last_qty = parse_betfair_quantity(fill_qty).ok()?;
447        let last_px = parse_betfair_price(fill_price).ok()?;
448
449        self.record_cumulative_state(bet_id, size_matched, average_price_matched);
450        self.published_trade_ids.insert(trade_id.to_string());
451        self.fill_lots
452            .entry(bet_id.to_string())
453            .or_default()
454            .push(FillLot {
455                trade_id,
456                quantity: fill_qty,
457                price: last_px,
458            });
459
460        Some((trade_id, last_qty, last_px))
461    }
462
463    /// Allocates cumulative Betfair `sv` to known applied fill lots, newest first.
464    ///
465    /// First-seen voids without known fill lots are recorded but emit no correction.
466    pub fn maybe_fill_voids(&mut self, uo: &UnmatchedOrder) -> Vec<FillVoidAllocation> {
467        let cumulative = normalize_betfair_quantity(uo.sv.unwrap_or(Decimal::ZERO));
468        let previous = self
469            .voided_qty
470            .get(&uo.id)
471            .copied()
472            .unwrap_or(Decimal::ZERO);
473
474        if cumulative <= previous {
475            return Vec::new();
476        }
477
478        let Some(lots) = self.fill_lots.get(&uo.id) else {
479            self.voided_qty.insert(uo.id.clone(), cumulative);
480            return Vec::new();
481        };
482        let mut remaining = cumulative - previous;
483        let mut desired = Vec::new();
484
485        for lot in lots.iter().rev() {
486            if remaining <= Decimal::ZERO {
487                break;
488            }
489            let key = (uo.id.clone(), lot.trade_id);
490            let prior = self.fill_voids.get(&key).copied().unwrap_or(Decimal::ZERO);
491            let available = lot.quantity.saturating_sub(prior);
492            let increment = remaining.min(available);
493            if increment > Decimal::ZERO {
494                desired.push((lot, prior + increment));
495                remaining -= increment;
496            }
497        }
498
499        if remaining > Decimal::ZERO {
500            log::warn!(
501                "Betfair cumulative void exceeds known fill lots for bet_id={}: sv={cumulative}, unmatched={remaining}",
502                uo.id,
503            );
504            return Vec::new();
505        }
506
507        let mut updates = Vec::new();
508
509        for (lot, allocation) in desired.into_iter().rev() {
510            let key = (uo.id.clone(), lot.trade_id);
511            let Ok(voided_qty) = parse_betfair_quantity(allocation) else {
512                continue;
513            };
514            self.fill_voids.insert(key, allocation);
515            updates.push(FillVoidAllocation {
516                trade_id: lot.trade_id,
517                voided_qty,
518                last_px: lot.price,
519            });
520        }
521        self.voided_qty.insert(uo.id.clone(), cumulative);
522        updates
523    }
524
525    #[must_use]
526    pub(crate) fn has_fill_lots(&self, bet_id: &str) -> bool {
527        self.fill_lots
528            .get(bet_id)
529            .is_some_and(|lots| !lots.is_empty())
530    }
531
532    /// Returns whether a cumulative Betfair void update has not yet been applied.
533    #[must_use]
534    pub fn has_unseen_fill_void(&self, uo: &UnmatchedOrder) -> bool {
535        let cumulative = normalize_betfair_quantity(uo.sv.unwrap_or(Decimal::ZERO));
536        let previous = self
537            .voided_qty
538            .get(&uo.id)
539            .copied()
540            .unwrap_or(Decimal::ZERO);
541        cumulative > previous
542    }
543
544    /// Returns whether a cumulative Betfair fill update has not yet been applied.
545    #[must_use]
546    pub fn has_unseen_fill(&self, uo: &UnmatchedOrder) -> bool {
547        let size_matched = normalize_betfair_quantity(uo.sm.unwrap_or(Decimal::ZERO));
548        let cumulative = if self.has_fill_lots(&uo.id) {
549            size_matched + normalize_betfair_quantity(uo.sv.unwrap_or(Decimal::ZERO))
550        } else {
551            size_matched
552        };
553        let previous = self
554            .filled_qty
555            .get(&uo.id)
556            .copied()
557            .unwrap_or(Decimal::ZERO);
558        let order_qty = normalize_betfair_quantity(resolve_stream_order_quantity(uo.s, uo));
559        cumulative > previous && cumulative <= order_qty
560    }
561
562    pub(crate) fn sync_fill_lot(
563        &mut self,
564        bet_id: &str,
565        trade_id: TradeId,
566        quantity: Decimal,
567        price: Price,
568        voided_qty: Decimal,
569    ) {
570        let lots = self.fill_lots.entry(bet_id.to_string()).or_default();
571        if !lots.iter().any(|lot| lot.trade_id == trade_id) {
572            lots.push(FillLot {
573                trade_id,
574                quantity: normalize_betfair_quantity(quantity),
575                price,
576            });
577        }
578        let gross_filled = lots.iter().map(|lot| lot.quantity).sum();
579        self.filled_qty.insert(bet_id.to_string(), gross_filled);
580        self.published_trade_ids.insert(trade_id.to_string());
581        if voided_qty > Decimal::ZERO {
582            self.fill_voids.insert(
583                (bet_id.to_string(), trade_id),
584                normalize_betfair_quantity(voided_qty),
585            );
586        }
587    }
588
589    pub(crate) fn matched_quantity(&self, bet_id: &str) -> Decimal {
590        let filled = self.filled_qty.get(bet_id).copied().unwrap_or_default();
591        let voided = self.voided_qty.get(bet_id).copied().unwrap_or_default();
592        (filled - voided).max(Decimal::ZERO)
593    }
594
595    pub(crate) fn sync_voided_qty(&mut self, bet_id: &str, voided_qty: Decimal) {
596        self.voided_qty
597            .insert(bet_id.to_string(), normalize_betfair_quantity(voided_qty));
598    }
599
600    fn record_cumulative_state(
601        &mut self,
602        bet_id: &str,
603        size_matched: Decimal,
604        average_price_matched: Option<Decimal>,
605    ) {
606        self.filled_qty.insert(bet_id.to_string(), size_matched);
607        if let Some(avg_px) = average_price_matched.map(normalize_betfair_price) {
608            self.avg_px.insert(bet_id.to_string(), avg_px);
609        }
610    }
611
612    /// Back-calculates the individual fill price from Betfair's cumulative
613    /// average price matched (`avp`).
614    ///
615    /// For the first fill, the average price IS the fill price. For subsequent
616    /// fills, the individual price is derived from:
617    /// `fill_price = (avp * sm - prev_avp * prev_sm) / fill_size`
618    fn compute_fill_price(
619        &self,
620        bet_id: &str,
621        average_price_matched: Option<Decimal>,
622        fallback_price: Decimal,
623        prev_filled: Decimal,
624        sm: Decimal,
625    ) -> Decimal {
626        let Some(avp) = average_price_matched.map(normalize_betfair_price) else {
627            return fallback_price;
628        };
629
630        if prev_filled == Decimal::ZERO {
631            return avp;
632        }
633
634        let Some(prev_avg) = self.avg_px.get(bet_id).copied() else {
635            return avp;
636        };
637
638        if prev_avg == avp {
639            return avp;
640        }
641
642        let fill_size = sm - prev_filled;
643
644        if fill_size == Decimal::ZERO {
645            return prev_avg;
646        }
647
648        let fill_price = (avp * sm - prev_avg * prev_filled) / fill_size;
649
650        if fill_price <= Decimal::ZERO {
651            log::warn!(
652                "Calculated fill price {fill_price} is invalid for bet_id={bet_id}, falling back to avp={avp}",
653            );
654            return avp;
655        }
656
657        fill_price
658    }
659
660    fn compute_fill_price_with_voids(
661        &self,
662        bet_id: &str,
663        average_price_matched: Option<Decimal>,
664        fallback_price: Decimal,
665        gross_matched: Decimal,
666        fill_qty: Decimal,
667        cumulative_voided: Decimal,
668    ) -> Decimal {
669        let Some(avp) = average_price_matched.map(normalize_betfair_price) else {
670            return fallback_price;
671        };
672        let cumulative_voided = normalize_betfair_quantity(cumulative_voided);
673        let previous_voided = self
674            .voided_qty
675            .get(bet_id)
676            .copied()
677            .map_or(Decimal::ZERO, normalize_betfair_quantity);
678        let incremental_void = cumulative_voided.saturating_sub(previous_voided);
679        let surviving_fill_qty = fill_qty.saturating_sub(incremental_void);
680        if surviving_fill_qty <= Decimal::ZERO {
681            return avp;
682        }
683
684        let mut prior_notional = Decimal::ZERO;
685        let mut effective_lots = Vec::new();
686
687        if let Some(lots) = self.fill_lots.get(bet_id) {
688            for lot in lots {
689                let already_voided = self
690                    .fill_voids
691                    .get(&(bet_id.to_string(), lot.trade_id))
692                    .copied()
693                    .unwrap_or(Decimal::ZERO);
694                let effective = lot.quantity.saturating_sub(already_voided);
695                prior_notional += effective * lot.price.as_decimal();
696                effective_lots.push((effective, lot.price.as_decimal()));
697            }
698        }
699
700        let mut prior_void = incremental_void.saturating_sub(fill_qty);
701        for (effective, price) in effective_lots.into_iter().rev() {
702            if prior_void <= Decimal::ZERO {
703                break;
704            }
705            let removed = prior_void.min(effective);
706            prior_notional -= removed * price;
707            prior_void -= removed;
708        }
709
710        let surviving_matched = gross_matched.saturating_sub(cumulative_voided);
711        let fill_price = (avp * surviving_matched - prior_notional) / surviving_fill_qty;
712        if fill_price <= Decimal::ZERO {
713            log::warn!(
714                "Calculated post-void fill price {fill_price} is invalid for bet_id={bet_id}, falling back to avp={avp}",
715            );
716            return avp;
717        }
718        fill_price
719    }
720
721    /// Pre-populates state for a bet from existing order data.
722    ///
723    /// Monotonic in `filled_qty`: keeps the larger of the cached and
724    /// in-memory values so a cache that lags the tracker (engine has
725    /// not yet processed an emitted fill) cannot regress it.
726    pub fn sync_order(&mut self, bet_id: &str, filled_qty: Decimal, avg_px: Decimal) {
727        let filled_qty = normalize_betfair_quantity(filled_qty);
728        let avg_px = normalize_betfair_price(avg_px);
729        let current = self
730            .filled_qty
731            .get(bet_id)
732            .copied()
733            .map_or(Decimal::ZERO, normalize_betfair_quantity);
734
735        if filled_qty > current {
736            self.filled_qty.insert(bet_id.to_string(), filled_qty);
737            if avg_px > Decimal::ZERO {
738                self.avg_px.insert(bet_id.to_string(), avg_px);
739            }
740        }
741    }
742
743    /// Seeds the trade-id dedup set so fills already published via another
744    /// channel are not re-emitted when the post-reconnect image arrives.
745    pub fn seed_published_trade_ids<I, S>(&mut self, trade_ids: I)
746    where
747        I: IntoIterator<Item = S>,
748        S: Into<String>,
749    {
750        for id in trade_ids {
751            self.published_trade_ids.insert(id.into());
752        }
753    }
754
755    /// Removes state for a completed bet to prevent unbounded growth.
756    pub fn prune(&mut self, bet_id: &str) {
757        self.filled_qty.remove(bet_id);
758        self.voided_qty.remove(bet_id);
759        self.avg_px.remove(bet_id);
760        self.fill_lots.remove(bet_id);
761        self.fill_voids.retain(|(id, _), _| id != bet_id);
762
763        let prefix = format!("{bet_id}-");
764        self.published_trade_ids
765            .retain(|id| !id.starts_with(&prefix));
766    }
767}
768
769fn make_trade_id_for_size(bet_id: &str, size_matched: Decimal) -> TradeId {
770    TradeId::new(format!("{bet_id}-{size_matched}"))
771}
772
773/// Returns `true` if the unmatched order has cancel, lapse, or void quantities.
774#[must_use]
775pub fn has_cancel_quantity(uo: &UnmatchedOrder) -> bool {
776    let sc = uo.sc.unwrap_or(Decimal::ZERO);
777    let sl = uo.sl.unwrap_or(Decimal::ZERO);
778    let sv = uo.sv.unwrap_or(Decimal::ZERO);
779    (sc + sl + sv) > Decimal::ZERO
780}
781
782/// Returns `true` if the order is execution-complete and has lapsed.
783#[must_use]
784pub fn is_lapsed(uo: &UnmatchedOrder) -> bool {
785    uo.status == StreamingOrderStatus::ExecutionComplete && uo.lsrc.is_some()
786}
787
788/// Returns `true` if the unmatched order is an SP (on-close) bet resting until
789/// BSP reconciliation: execution-complete with no matched, cancelled, lapsed,
790/// or voided quantity and no lapse reason.
791///
792/// SP bets cannot be cancelled once placed, so this state is open, not
793/// terminal.
794#[must_use]
795pub fn is_resting_sp_bet(uo: &UnmatchedOrder) -> bool {
796    uo.status == StreamingOrderStatus::ExecutionComplete
797        && matches!(
798            uo.ot,
799            StreamingOrderType::LimitOnClose | StreamingOrderType::MarketOnClose
800        )
801        && uo.sm.unwrap_or(Decimal::ZERO) <= Decimal::ZERO
802        && uo.lsrc.is_none()
803        && !has_cancel_quantity(uo)
804}
805
806/// Parses a streaming [`UnmatchedOrder`] into a Nautilus [`OrderStatusReport`].
807///
808/// Resolves the Nautilus order status from the Betfair streaming status
809/// plus matched/cancelled quantities.
810///
811/// # Errors
812///
813/// Returns an error if price or quantity values cannot be converted.
814pub fn parse_order_status_report(
815    uo: &UnmatchedOrder,
816    instrument_id: InstrumentId,
817    account_id: AccountId,
818    ts_event: UnixNanos,
819    ts_init: UnixNanos,
820) -> anyhow::Result<OrderStatusReport> {
821    let order_side = OrderSide::from(uo.side);
822    let order_type = OrderType::from(uo.ot);
823    let time_in_force = parse_stream_time_in_force(uo)?;
824
825    let size_matched = uo.sm.unwrap_or(Decimal::ZERO);
826    let size_cancelled = uo.sc.unwrap_or(Decimal::ZERO);
827    let size_lapsed = uo.sl.unwrap_or(Decimal::ZERO);
828    let size_voided = uo.sv.unwrap_or(Decimal::ZERO);
829
830    // Include lapsed/voided in the closed quantity for status resolution
831    let size_closed = size_cancelled + size_lapsed + size_voided;
832    let order_status = if uo.status == StreamingOrderStatus::ExecutionComplete
833        && size_voided > Decimal::ZERO
834        && size_cancelled.is_zero()
835        && size_lapsed.is_zero()
836    {
837        OrderStatus::Voided
838    } else if is_resting_sp_bet(uo) {
839        OrderStatus::Accepted
840    } else {
841        resolve_streaming_order_status(uo.status, size_matched, size_closed)
842    };
843
844    let quantity_decimal = stream_order_quantity(uo);
845    anyhow::ensure!(
846        quantity_decimal > Decimal::ZERO,
847        "failed to resolve positive quantity for stream order update {} \
848         (order_type={:?}, persistence_type={:?}, size={}, bsp_liability={:?}, \
849         size_matched={:?}, size_remaining={:?}, size_cancelled={:?}, size_lapsed={:?}, size_voided={:?})",
850        uo.id,
851        uo.ot,
852        uo.pt,
853        uo.s,
854        uo.bsp,
855        uo.sm,
856        uo.sr,
857        uo.sc,
858        uo.sl,
859        uo.sv,
860    );
861    let quantity = parse_betfair_quantity(quantity_decimal)?;
862    let filled_qty = parse_betfair_quantity(size_matched)?;
863
864    let ts_accepted = parse_millis_timestamp(uo.pd);
865
866    // Use the latest lifecycle timestamp, falling back to OCM publish time
867    let ts_last = [uo.md, uo.cd, uo.ld]
868        .into_iter()
869        .flatten()
870        .max()
871        .map_or(ts_event, parse_millis_timestamp);
872
873    let venue_order_id = VenueOrderId::from(uo.id.as_str());
874    let client_order_id = uo
875        .rfo
876        .as_deref()
877        .filter(|s| !s.is_empty())
878        .map(ClientOrderId::from);
879
880    let price = parse_betfair_price(uo.p)?;
881
882    let mut report = OrderStatusReport::new(
883        account_id,
884        instrument_id,
885        client_order_id,
886        venue_order_id,
887        order_side.into(),
888        order_type,
889        time_in_force,
890        order_status,
891        quantity,
892        filled_qty,
893        ts_accepted,
894        ts_last,
895        ts_init,
896        None,
897    )
898    .with_price(price);
899
900    report.avg_px = uo.avp;
901    if let Some(lsrc) = uo.lsrc {
902        report.cancel_reason = Some(lsrc.to_string());
903    }
904
905    Ok(report)
906}
907
908fn parse_stream_time_in_force(uo: &UnmatchedOrder) -> anyhow::Result<TimeInForce> {
909    if matches!(
910        uo.ot,
911        StreamingOrderType::LimitOnClose | StreamingOrderType::MarketOnClose
912    ) {
913        // The stream can carry a non-BSP persistence type on SP bets; the
914        // on-close instruction defines the time in force
915        return Ok(TimeInForce::AtTheClose);
916    }
917
918    match uo.pt {
919        Some(persistence_type) => Ok(TimeInForce::from(persistence_type)),
920        None => anyhow::bail!("missing persistence type for order update {}", uo.id),
921    }
922}
923
924fn stream_order_quantity(uo: &UnmatchedOrder) -> Decimal {
925    if uo.s > Decimal::ZERO {
926        return uo.s;
927    }
928
929    let lifecycle_qty = uo.sm.unwrap_or(Decimal::ZERO)
930        + uo.sr.unwrap_or(Decimal::ZERO)
931        + uo.sc.unwrap_or(Decimal::ZERO)
932        + uo.sl.unwrap_or(Decimal::ZERO)
933        + uo.sv.unwrap_or(Decimal::ZERO);
934
935    if lifecycle_qty > Decimal::ZERO {
936        return lifecycle_qty;
937    }
938
939    if uses_liability_based_stream_quantity(uo) {
940        return uo.bsp.unwrap_or(Decimal::ZERO);
941    }
942
943    Decimal::ZERO
944}
945
946fn resolve_stream_order_quantity(order_qty: Decimal, uo: &UnmatchedOrder) -> Decimal {
947    if order_qty > Decimal::ZERO {
948        order_qty
949    } else {
950        stream_order_quantity(uo)
951    }
952}
953
954fn uses_liability_based_stream_quantity(uo: &UnmatchedOrder) -> bool {
955    matches!(
956        uo.ot,
957        crate::common::enums::StreamingOrderType::LimitOnClose
958            | crate::common::enums::StreamingOrderType::MarketOnClose
959    )
960}
961
962/// Creates a [`FillReport`] for a Betfair order fill.
963///
964/// Betfair charges commission on net winnings, not per-fill, so commission
965/// is set to zero. The `liquidity_side` is unknown from the stream.
966#[must_use]
967#[expect(clippy::too_many_arguments)]
968pub fn make_fill_report(
969    account_id: AccountId,
970    instrument_id: InstrumentId,
971    venue_order_id: VenueOrderId,
972    trade_id: TradeId,
973    order_side: OrderSide,
974    last_qty: Quantity,
975    last_px: Price,
976    currency: Currency,
977    client_order_id: Option<ClientOrderId>,
978    ts_event: UnixNanos,
979    ts_init: UnixNanos,
980) -> FillReport {
981    FillReport::new(
982        account_id,
983        instrument_id,
984        venue_order_id,
985        trade_id,
986        order_side,
987        last_qty,
988        last_px,
989        Money::zero(currency),
990        LiquiditySide::NoLiquiditySide,
991        client_order_id,
992        None,
993        ts_event,
994        ts_init,
995        None,
996    )
997}
998
999/// Extracts a [`BetfairTicker`] from a runner change if any ticker fields are present.
1000///
1001/// Returns `None` when the runner change contains no ltp, tv, spn, or spf data.
1002#[must_use]
1003pub fn parse_betfair_ticker(
1004    instrument_id: InstrumentId,
1005    rc: &RunnerChange,
1006    ts_event: UnixNanos,
1007    ts_init: UnixNanos,
1008) -> Option<BetfairTicker> {
1009    if rc.ltp.is_none() && rc.tv.is_none() && rc.spn.is_none() && rc.spf.is_none() {
1010        return None;
1011    }
1012
1013    Some(BetfairTicker::new(
1014        instrument_id,
1015        rc.ltp,
1016        rc.tv,
1017        rc.spn,
1018        rc.spf,
1019        ts_event,
1020        ts_init,
1021    ))
1022}
1023
1024/// Extracts [`BetfairStartingPrice`] values from a market definition's runners.
1025///
1026/// Returns one entry per runner that has a non-None BSP value.
1027#[must_use]
1028pub fn parse_betfair_starting_prices(
1029    market_id: &str,
1030    def: &MarketDefinition,
1031    ts_event: UnixNanos,
1032    ts_init: UnixNanos,
1033) -> Vec<BetfairStartingPrice> {
1034    let Some(runners) = &def.runners else {
1035        return Vec::new();
1036    };
1037
1038    runners
1039        .iter()
1040        .filter_map(|rd| {
1041            let bsp = rd.bsp?;
1042            let handicap = rd.hc.unwrap_or(Decimal::ZERO);
1043            let instrument_id = make_instrument_id(market_id, rd.id, handicap);
1044            Some(BetfairStartingPrice::new(
1045                instrument_id,
1046                bsp,
1047                ts_event,
1048                ts_init,
1049            ))
1050        })
1051        .collect()
1052}
1053
1054/// Extracts BSP order book deltas from a runner change's `spb`/`spl` fields.
1055///
1056/// Returns an empty vec when neither `spb` nor `spl` data is present.
1057/// The `side` field uses `OrderSide::Sell` for `spb` (back) and
1058/// `OrderSide::Buy` for `spl` (lay), following Betfair's inverted convention.
1059#[must_use]
1060pub fn parse_bsp_book_deltas(
1061    instrument_id: InstrumentId,
1062    rc: &RunnerChange,
1063    ts_event: UnixNanos,
1064    ts_init: UnixNanos,
1065) -> Vec<BetfairBspBookDelta> {
1066    let spb_len = rc.spb.as_ref().map_or(0, Vec::len);
1067    let spl_len = rc.spl.as_ref().map_or(0, Vec::len);
1068
1069    if spb_len + spl_len == 0 {
1070        return Vec::new();
1071    }
1072
1073    let mut result = Vec::with_capacity(spb_len + spl_len);
1074
1075    for (levels, side) in [
1076        (rc.spb.as_deref().unwrap_or(&[]), OrderSide::Sell),
1077        (rc.spl.as_deref().unwrap_or(&[]), OrderSide::Buy),
1078    ] {
1079        for pv in levels {
1080            let action = if pv.volume == Decimal::ZERO {
1081                BookAction::Delete
1082            } else {
1083                BookAction::Update
1084            };
1085
1086            result.push(BetfairBspBookDelta::new(
1087                instrument_id,
1088                action,
1089                side,
1090                pv.price,
1091                pv.volume,
1092                ts_event,
1093                ts_init,
1094            ));
1095        }
1096    }
1097
1098    result
1099}
1100
1101/// Produces [`InstrumentClose`] events from a market definition's runner statuses.
1102///
1103/// Winners and placed runners get close price 1.0; losers and removed runners
1104/// get close price 0.0. Active runners produce no close event.
1105#[must_use]
1106pub fn parse_instrument_closes(
1107    market_id: &str,
1108    def: &MarketDefinition,
1109    ts_event: UnixNanos,
1110    ts_init: UnixNanos,
1111) -> Vec<InstrumentClose> {
1112    let Some(runners) = &def.runners else {
1113        return Vec::new();
1114    };
1115
1116    runners
1117        .iter()
1118        .filter_map(|rd| {
1119            let status = rd.status.as_ref()?;
1120            let close_price = match status {
1121                RunnerStatus::Winner | RunnerStatus::Placed => Price::from("1.00"),
1122                RunnerStatus::Loser | RunnerStatus::Removed | RunnerStatus::RemovedVacant => {
1123                    Price::from("0.00")
1124                }
1125                RunnerStatus::Active | RunnerStatus::Hidden | RunnerStatus::Unknown => return None,
1126            };
1127
1128            let handicap = rd.hc.unwrap_or(Decimal::ZERO);
1129            let instrument_id = make_instrument_id(market_id, rd.id, handicap);
1130
1131            Some(InstrumentClose::new(
1132                instrument_id,
1133                close_price,
1134                InstrumentCloseType::ContractExpired,
1135                ts_event,
1136                ts_init,
1137            ))
1138        })
1139        .collect()
1140}
1141
1142/// Parses a single [`RaceRunnerChange`] into a [`BetfairRaceRunnerData`].
1143///
1144/// Returns `None` if the runner change has no selection ID.
1145#[must_use]
1146pub fn parse_race_runner_data(
1147    race_id: &str,
1148    market_id: &str,
1149    rrc: &RaceRunnerChange,
1150    ts_event: UnixNanos,
1151    ts_init: UnixNanos,
1152) -> Option<BetfairRaceRunnerData> {
1153    let selection_id = rrc.id?;
1154
1155    Some(BetfairRaceRunnerData::new(
1156        race_id.to_string(),
1157        market_id.to_string(),
1158        selection_id,
1159        rrc.lat.unwrap_or(f64::NAN),
1160        rrc.lng.unwrap_or(f64::NAN),
1161        rrc.spd.unwrap_or(f64::NAN),
1162        rrc.prg.unwrap_or(f64::NAN),
1163        rrc.sfq.unwrap_or(f64::NAN),
1164        ts_event,
1165        ts_init,
1166    ))
1167}
1168
1169/// Parses a [`RaceProgressChange`] into a [`BetfairRaceProgress`].
1170#[must_use]
1171pub fn parse_race_progress(
1172    race_id: &str,
1173    market_id: &str,
1174    rpc: &RaceProgressChange,
1175    ts_event: UnixNanos,
1176    ts_init: UnixNanos,
1177) -> BetfairRaceProgress {
1178    let order_json = serialize_json_value(rpc.ord.as_ref());
1179    let jumps_json = serialize_json_value(rpc.jumps.as_ref());
1180
1181    BetfairRaceProgress::new(
1182        race_id.to_string(),
1183        market_id.to_string(),
1184        rpc.g.clone().unwrap_or_default(),
1185        rpc.st.unwrap_or(f64::NAN),
1186        rpc.rt.unwrap_or(f64::NAN),
1187        rpc.spd.unwrap_or(f64::NAN),
1188        rpc.prg.unwrap_or(f64::NAN),
1189        order_json,
1190        jumps_json,
1191        ts_event,
1192        ts_init,
1193    )
1194}
1195
1196#[must_use]
1197pub fn parse_cricket_match(
1198    cricket: &CricketChange,
1199    ts_event: UnixNanos,
1200    ts_init: UnixNanos,
1201) -> Option<BetfairCricketMatch> {
1202    let event_id = cricket.event_id.clone()?;
1203    let market_id = cricket.market_id.clone()?;
1204
1205    Some(BetfairCricketMatch::new(
1206        event_id,
1207        market_id,
1208        serialize_json_value(cricket.fixture_info.as_ref()),
1209        serialize_json_value(cricket.home_team.as_ref()),
1210        serialize_json_value(cricket.away_team.as_ref()),
1211        serialize_json_value(cricket.match_stats.as_ref()),
1212        serialize_json_value(cricket.incident_list_wrapper.as_ref()),
1213        ts_event,
1214        ts_init,
1215    ))
1216}
1217
1218fn serialize_json_value<T>(value: Option<&T>) -> String
1219where
1220    T: serde::Serialize,
1221{
1222    value
1223        .map(|value| serde_json::to_string(value).unwrap_or_default())
1224        .unwrap_or_default()
1225}
1226
1227#[cfg(test)]
1228mod tests {
1229    use nautilus_model::enums::{MarketStatusAction, OrderStatus, TimeInForce};
1230    use rstest::rstest;
1231
1232    use super::*;
1233    use crate::{
1234        common::{
1235            consts::{BETFAIR_PRICE_PRECISION, BETFAIR_QUANTITY_PRECISION},
1236            enums::{StreamingOrderType, StreamingPersistenceType, StreamingSide},
1237            testing::load_test_json,
1238        },
1239        stream::messages::{PV, RunnerChange, RunnerDefinition, StreamMessage, stream_decode},
1240    };
1241
1242    fn assert_decimal_option_eq(actual: Option<Decimal>, expected: Decimal) {
1243        assert_eq!(actual, Some(expected));
1244    }
1245
1246    #[rstest]
1247    fn test_parse_runner_book_snapshot() {
1248        let data = load_test_json("stream/mcm_live_IMAGE.json");
1249        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1250
1251        if let StreamMessage::MarketChange(mcm) = msg {
1252            let mc = mcm.mc.as_ref().unwrap();
1253            let change = &mc[0];
1254            let rc = &change.rc.as_ref().unwrap()[0];
1255            let instrument_id = make_instrument_id(&change.id, rc.id, Decimal::ZERO);
1256
1257            let deltas = parse_runner_book_deltas(
1258                instrument_id,
1259                rc,
1260                true,
1261                mcm.pt,
1262                parse_millis_timestamp(mcm.pt),
1263                parse_millis_timestamp(mcm.pt),
1264            )
1265            .unwrap()
1266            .expect("should produce deltas");
1267
1268            // Clear + atb levels + atl levels
1269            let atb_len = rc.atb.as_ref().unwrap().len();
1270            let atl_len = rc.atl.as_ref().unwrap().len();
1271            assert_eq!(deltas.deltas.len(), 1 + atb_len + atl_len);
1272
1273            // First delta is Clear with F_SNAPSHOT
1274            assert_eq!(deltas.deltas[0].action, BookAction::Clear);
1275            assert!(RecordFlag::F_SNAPSHOT.matches(deltas.deltas[0].flags));
1276
1277            // Subsequent deltas are Add with F_SNAPSHOT
1278            for delta in &deltas.deltas[1..] {
1279                assert_eq!(delta.action, BookAction::Add);
1280                assert!(RecordFlag::F_SNAPSHOT.matches(delta.flags));
1281            }
1282
1283            // Last delta has F_LAST
1284            let last = deltas.deltas.last().unwrap();
1285            assert!(RecordFlag::F_LAST.matches(last.flags));
1286
1287            // Verify buy/sell sides
1288            let buy_count = deltas
1289                .deltas
1290                .iter()
1291                .filter(|d| d.order.side == Some(OrderSide::Buy))
1292                .count();
1293            let sell_count = deltas
1294                .deltas
1295                .iter()
1296                .filter(|d| d.order.side == Some(OrderSide::Sell))
1297                .count();
1298            assert_eq!(buy_count, atb_len);
1299            assert_eq!(sell_count, atl_len);
1300        } else {
1301            panic!("expected MarketChange");
1302        }
1303    }
1304
1305    #[rstest]
1306    fn test_parse_runner_book_snapshot_skips_zero_volume_levels() {
1307        let rc: RunnerChange = serde_json::from_str(
1308            r#"{
1309                "id": 123,
1310                "atb": [[2.0, 0.0], [2.1, 3.0]],
1311                "atl": [[2.2, 0.0], [2.3, 4.0]]
1312            }"#,
1313        )
1314        .unwrap();
1315        let instrument_id = make_instrument_id("1.234", rc.id, Decimal::ZERO);
1316
1317        let deltas = parse_runner_book_deltas(
1318            instrument_id,
1319            &rc,
1320            true,
1321            1_551_400_000_000,
1322            parse_millis_timestamp(1_551_400_000_000),
1323            parse_millis_timestamp(1_551_400_000_000),
1324        )
1325        .unwrap()
1326        .expect("should produce deltas");
1327
1328        assert_eq!(deltas.deltas.len(), 3);
1329        assert_eq!(deltas.deltas[0].action, BookAction::Clear);
1330        assert!(
1331            deltas
1332                .deltas
1333                .iter()
1334                .filter(|delta| delta.action == BookAction::Add)
1335                .all(|delta| !delta.order.size.is_zero())
1336        );
1337    }
1338
1339    #[rstest]
1340    fn test_parse_runner_book_update() {
1341        let data = load_test_json("stream/mcm_UPDATE.json");
1342        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1343
1344        if let StreamMessage::MarketChange(mcm) = msg {
1345            let mc = mcm.mc.as_ref().unwrap();
1346            let change = &mc[0];
1347            let rc = &change.rc.as_ref().unwrap()[0];
1348            let instrument_id = make_instrument_id(&change.id, rc.id, Decimal::ZERO);
1349
1350            let deltas = parse_runner_book_deltas(
1351                instrument_id,
1352                rc,
1353                false,
1354                mcm.pt,
1355                parse_millis_timestamp(mcm.pt),
1356                parse_millis_timestamp(mcm.pt),
1357            )
1358            .unwrap()
1359            .expect("should produce deltas");
1360
1361            // No Clear delta for updates
1362            assert!(deltas.deltas.iter().all(|d| d.action != BookAction::Clear));
1363
1364            // Last delta has F_LAST
1365            let last = deltas.deltas.last().unwrap();
1366            assert!(RecordFlag::F_LAST.matches(last.flags));
1367
1368            // No snapshot flags
1369            for delta in &deltas.deltas {
1370                assert!(!RecordFlag::F_SNAPSHOT.matches(delta.flags));
1371            }
1372        } else {
1373            panic!("expected MarketChange");
1374        }
1375    }
1376
1377    #[rstest]
1378    fn test_parse_runner_book_update_zero_volume_is_delete() {
1379        let data = load_test_json("stream/mcm_UPDATE.json");
1380        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1381
1382        if let StreamMessage::MarketChange(mcm) = msg {
1383            let mc = mcm.mc.as_ref().unwrap();
1384            let change = &mc[0];
1385            let rc = &change.rc.as_ref().unwrap()[0];
1386            let instrument_id = make_instrument_id(&change.id, rc.id, Decimal::ZERO);
1387
1388            let deltas = parse_runner_book_deltas(
1389                instrument_id,
1390                rc,
1391                false,
1392                mcm.pt,
1393                parse_millis_timestamp(mcm.pt),
1394                parse_millis_timestamp(mcm.pt),
1395            )
1396            .unwrap()
1397            .unwrap();
1398
1399            // atl has [[4.7, 0]] which should be Delete
1400            assert!(
1401                deltas.deltas.iter().any(|d| d.action == BookAction::Delete),
1402                "zero volume should produce Delete action"
1403            );
1404        } else {
1405            panic!("expected MarketChange");
1406        }
1407    }
1408
1409    #[rstest]
1410    fn test_parse_runner_book_no_levels_returns_none() {
1411        let rc = RunnerChange {
1412            id: 12345,
1413            hc: None,
1414            atb: None,
1415            atl: None,
1416            batb: None,
1417            batl: None,
1418            bdatb: None,
1419            bdatl: None,
1420            spb: None,
1421            spl: None,
1422            spn: None,
1423            spf: None,
1424            trd: None,
1425            ltp: None,
1426            tv: None,
1427        };
1428
1429        let result = parse_runner_book_deltas(
1430            make_instrument_id("1.123", 12345, Decimal::ZERO),
1431            &rc,
1432            false,
1433            0,
1434            UnixNanos::default(),
1435            UnixNanos::default(),
1436        );
1437
1438        assert!(result.unwrap().is_none());
1439    }
1440
1441    #[rstest]
1442    fn test_make_trade_tick() {
1443        let instrument_id = make_instrument_id("1.180737206", 19248890, Decimal::ZERO);
1444        let tick = make_trade_tick(
1445            instrument_id,
1446            Price::new(2.42, BETFAIR_PRICE_PRECISION),
1447            Quantity::new(100.0, BETFAIR_QUANTITY_PRECISION),
1448            TradeId::from("test-trade-1"),
1449            UnixNanos::default(),
1450            UnixNanos::default(),
1451        );
1452
1453        assert_eq!(tick.instrument_id, instrument_id);
1454        assert_eq!(tick.price.as_f64(), 2.42);
1455        assert_eq!(tick.size.as_f64(), 100.0);
1456        assert_eq!(tick.aggressor_side, AggressorSide::NoAggressor);
1457    }
1458
1459    fn make_status_def(
1460        status: MarketStatus,
1461        in_play: bool,
1462        runner_status: RunnerStatus,
1463    ) -> MarketDefinition {
1464        MarketDefinition {
1465            runners: Some(vec![RunnerDefinition {
1466                id: 456,
1467                hc: None,
1468                sort_priority: None,
1469                name: None,
1470                status: Some(runner_status),
1471                adjustment_factor: None,
1472                bsp: None,
1473                removal_date: None,
1474            }]),
1475            bet_delay: None,
1476            betting_type: None,
1477            bsp_market: None,
1478            bsp_reconciled: None,
1479            competition_id: None,
1480            competition_name: None,
1481            complete: None,
1482            country_code: None,
1483            cross_matching: None,
1484            discount_allowed: None,
1485            each_way_divisor: None,
1486            event_id: None,
1487            event_name: None,
1488            event_type_id: None,
1489            event_type_name: None,
1490            in_play: Some(in_play),
1491            line_interval: None,
1492            line_max_unit: None,
1493            line_min_unit: None,
1494            market_base_rate: None,
1495            market_id: None,
1496            market_name: None,
1497            market_time: None,
1498            market_type: None,
1499            number_of_active_runners: None,
1500            number_of_winners: None,
1501            open_date: None,
1502            persistence_enabled: None,
1503            price_ladder_definition: None,
1504            race_type: None,
1505            regulators: None,
1506            runners_voidable: None,
1507            settled_time: None,
1508            status: Some(status),
1509            suspend_time: None,
1510            timezone: None,
1511            turn_in_play_enabled: None,
1512            venue: None,
1513            version: None,
1514        }
1515    }
1516
1517    #[rstest]
1518    #[case(MarketStatus::Open, false, MarketStatusAction::PreOpen, false)]
1519    #[case(MarketStatus::Open, true, MarketStatusAction::Trading, true)]
1520    #[case(MarketStatus::Closed, false, MarketStatusAction::Close, false)]
1521    #[case(MarketStatus::Closed, true, MarketStatusAction::Close, false)]
1522    #[case(MarketStatus::Suspended, false, MarketStatusAction::Pause, false)]
1523    #[case(MarketStatus::Suspended, true, MarketStatusAction::Pause, false)]
1524    #[case(MarketStatus::Inactive, false, MarketStatusAction::Close, false)]
1525    fn test_parse_instrument_statuses_market_state(
1526        #[case] status: MarketStatus,
1527        #[case] in_play: bool,
1528        #[case] expected_action: MarketStatusAction,
1529        #[case] expected_is_trading: bool,
1530    ) {
1531        let def = make_status_def(status, in_play, RunnerStatus::Active);
1532        let results =
1533            parse_instrument_statuses("1.123", &def, UnixNanos::default(), UnixNanos::default());
1534
1535        assert_eq!(results.len(), 1);
1536        assert_eq!(results[0].action, expected_action);
1537        assert_eq!(results[0].is_trading, Some(expected_is_trading));
1538    }
1539
1540    #[rstest]
1541    fn test_parse_instrument_statuses_skips_unknown_market_status() {
1542        // An unmodeled market status must not publish a fabricated non-trading
1543        // state; the whole definition is skipped.
1544        let def = make_status_def(MarketStatus::Unknown, true, RunnerStatus::Active);
1545        let results =
1546            parse_instrument_statuses("1.123", &def, UnixNanos::default(), UnixNanos::default());
1547
1548        assert!(results.is_empty());
1549    }
1550
1551    #[rstest]
1552    fn test_parse_instrument_statuses_skips_unknown_runner() {
1553        // An unmodeled runner status must not be emitted as tradable, even in an
1554        // open in-play market; we skip it rather than fabricate tradability.
1555        let def = make_status_def(MarketStatus::Open, true, RunnerStatus::Unknown);
1556        let results =
1557            parse_instrument_statuses("1.123", &def, UnixNanos::default(), UnixNanos::default());
1558
1559        assert!(results.is_empty());
1560    }
1561
1562    #[rstest]
1563    #[case(RunnerStatus::Removed)]
1564    #[case(RunnerStatus::RemovedVacant)]
1565    fn test_parse_instrument_statuses_scratched_runner_closes(#[case] runner_status: RunnerStatus) {
1566        // Even with Open + in_play the runner must close when scratched
1567        let def = make_status_def(MarketStatus::Open, true, runner_status);
1568        let results =
1569            parse_instrument_statuses("1.123", &def, UnixNanos::default(), UnixNanos::default());
1570
1571        assert_eq!(results.len(), 1);
1572        assert_eq!(results[0].action, MarketStatusAction::Close);
1573        assert_eq!(results[0].is_trading, Some(false));
1574    }
1575
1576    #[rstest]
1577    #[case::missing_runners("runners")]
1578    #[case::missing_status("status")]
1579    fn test_parse_instrument_statuses_returns_empty(#[case] drop_field: &str) {
1580        let mut def = make_status_def(MarketStatus::Open, true, RunnerStatus::Active);
1581        match drop_field {
1582            "runners" => def.runners = None,
1583            "status" => def.status = None,
1584            _ => unreachable!(),
1585        }
1586
1587        let results =
1588            parse_instrument_statuses("1.123", &def, UnixNanos::default(), UnixNanos::default());
1589
1590        assert!(results.is_empty());
1591    }
1592
1593    #[rstest]
1594    fn test_parse_instrument_statuses_mixed_runners() {
1595        // Three runners in one market definition: an Active runner should follow the
1596        // market-level mapping while Removed/RemovedVacant override to Close. Verify
1597        // that selection id and handicap propagate into the emitted instrument id.
1598        let mut def = make_status_def(MarketStatus::Open, true, RunnerStatus::Active);
1599        def.runners = Some(vec![
1600            RunnerDefinition {
1601                id: 101,
1602                hc: None,
1603                sort_priority: Some(1),
1604                name: None,
1605                status: Some(RunnerStatus::Active),
1606                adjustment_factor: None,
1607                bsp: None,
1608                removal_date: None,
1609            },
1610            RunnerDefinition {
1611                id: 202,
1612                hc: Some(Decimal::new(25, 1)), // 2.5 handicap
1613                sort_priority: Some(2),
1614                name: None,
1615                status: Some(RunnerStatus::Removed),
1616                adjustment_factor: None,
1617                bsp: None,
1618                removal_date: None,
1619            },
1620            RunnerDefinition {
1621                id: 303,
1622                hc: None,
1623                sort_priority: Some(3),
1624                name: None,
1625                status: Some(RunnerStatus::RemovedVacant),
1626                adjustment_factor: None,
1627                bsp: None,
1628                removal_date: None,
1629            },
1630        ]);
1631
1632        let results =
1633            parse_instrument_statuses("1.999", &def, UnixNanos::default(), UnixNanos::default());
1634
1635        assert_eq!(results.len(), 3);
1636
1637        assert_eq!(results[0].action, MarketStatusAction::Trading);
1638        assert_eq!(results[0].is_trading, Some(true));
1639
1640        assert_eq!(results[1].action, MarketStatusAction::Close);
1641        assert_eq!(results[1].is_trading, Some(false));
1642
1643        assert_eq!(results[2].action, MarketStatusAction::Close);
1644        assert_eq!(results[2].is_trading, Some(false));
1645
1646        // Each runner produces a distinct instrument id (selection id + handicap)
1647        assert_ne!(results[0].instrument_id, results[1].instrument_id);
1648        assert_ne!(results[1].instrument_id, results[2].instrument_id);
1649        assert_ne!(results[0].instrument_id, results[2].instrument_id);
1650    }
1651
1652    #[rstest]
1653    fn test_parse_order_status_report_new_order() {
1654        let data = load_test_json("stream/ocm_NEW_FULL_IMAGE.json");
1655        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1656
1657        if let StreamMessage::OrderChange(ocm) = msg {
1658            let oc = ocm.oc.as_ref().unwrap();
1659            let omc = &oc[0];
1660            let orc = &omc.orc.as_ref().unwrap()[0];
1661            let uo = &orc.uo.as_ref().unwrap()[0];
1662
1663            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
1664            let report = parse_order_status_report(
1665                uo,
1666                instrument_id,
1667                AccountId::from("BETFAIR-001"),
1668                parse_millis_timestamp(ocm.pt),
1669                parse_millis_timestamp(ocm.pt),
1670            )
1671            .unwrap();
1672
1673            // Partially filled: sm=4.75, sr=0.25, status=E
1674            assert_eq!(report.order_status, OrderStatus::PartiallyFilled);
1675            assert_eq!(report.order_side, Some(OrderSide::Sell)); // Back → Sell
1676            assert_eq!(report.order_type, OrderType::Limit);
1677            assert_eq!(report.filled_qty.as_f64(), 4.75);
1678            assert_eq!(report.quantity.as_f64(), 5.0);
1679            assert!(report.price.is_some());
1680            assert_eq!(report.price.unwrap().as_f64(), 12.0);
1681        } else {
1682            panic!("expected OrderChange");
1683        }
1684    }
1685
1686    #[rstest]
1687    fn test_parse_order_status_report_filled() {
1688        let data = load_test_json("stream/ocm_FILLED.json");
1689        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1690
1691        if let StreamMessage::OrderChange(ocm) = msg {
1692            let oc = ocm.oc.as_ref().unwrap();
1693            let omc = &oc[0];
1694            let orc = &omc.orc.as_ref().unwrap()[0];
1695            let uo = &orc.uo.as_ref().unwrap()[0];
1696
1697            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
1698            let report = parse_order_status_report(
1699                uo,
1700                instrument_id,
1701                AccountId::from("BETFAIR-001"),
1702                parse_millis_timestamp(ocm.pt),
1703                parse_millis_timestamp(ocm.pt),
1704            )
1705            .unwrap();
1706
1707            assert_eq!(report.order_status, OrderStatus::Filled);
1708            assert_eq!(report.order_side, Some(OrderSide::Buy)); // Lay → Buy
1709            assert_eq!(report.filled_qty.as_f64(), 10.0);
1710            assert_eq!(report.quantity.as_f64(), 10.0);
1711
1712            // Has client_order_id from rfo field
1713            assert!(report.client_order_id.is_some());
1714        } else {
1715            panic!("expected OrderChange");
1716        }
1717    }
1718
1719    #[rstest]
1720    fn test_parse_order_status_report_cancelled() {
1721        let data = load_test_json("stream/ocm_CANCEL.json");
1722        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1723
1724        if let StreamMessage::OrderChange(ocm) = msg {
1725            let oc = ocm.oc.as_ref().unwrap();
1726            let omc = &oc[0];
1727            let orc = &omc.orc.as_ref().unwrap()[0];
1728            let uo = &orc.uo.as_ref().unwrap()[0];
1729
1730            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
1731            let report = parse_order_status_report(
1732                uo,
1733                instrument_id,
1734                AccountId::from("BETFAIR-001"),
1735                parse_millis_timestamp(ocm.pt),
1736                parse_millis_timestamp(ocm.pt),
1737            )
1738            .unwrap();
1739
1740            assert_eq!(report.order_status, OrderStatus::Canceled);
1741            assert_eq!(report.order_side, Some(OrderSide::Sell)); // Back → Sell
1742            assert_eq!(report.filled_qty.as_f64(), 0.0);
1743            assert_eq!(report.quantity.as_f64(), 10.0);
1744        } else {
1745            panic!("expected OrderChange");
1746        }
1747    }
1748
1749    #[rstest]
1750    fn test_parse_order_status_report_duplicate_execution() {
1751        let data = load_test_json("stream/ocm_DUPLICATE_EXECUTION.json");
1752        let msgs: Vec<StreamMessage> = serde_json::from_str(&data).unwrap();
1753
1754        if let StreamMessage::OrderChange(ocm) = &msgs[0] {
1755            let oc = ocm.oc.as_ref().unwrap();
1756            let omc = &oc[0];
1757            let orc = &omc.orc.as_ref().unwrap()[0];
1758            let uo = &orc.uo.as_ref().unwrap()[0];
1759
1760            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
1761            let report = parse_order_status_report(
1762                uo,
1763                instrument_id,
1764                AccountId::from("BETFAIR-001"),
1765                parse_millis_timestamp(ocm.pt),
1766                parse_millis_timestamp(ocm.pt),
1767            )
1768            .unwrap();
1769
1770            // Partially filled: sm=1.12, status=E
1771            assert_eq!(report.order_status, OrderStatus::PartiallyFilled);
1772            assert_eq!(report.order_side, Some(OrderSide::Buy)); // Lay → Buy
1773            assert_eq!(report.filled_qty.as_f64(), 1.12);
1774        } else {
1775            panic!("expected OrderChange");
1776        }
1777    }
1778
1779    #[rstest]
1780    fn test_parse_order_status_report_sp_resting() {
1781        let data = load_test_json("stream/ocm_SP_RESTING.json");
1782        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1783
1784        if let StreamMessage::OrderChange(ocm) = msg {
1785            let oc = ocm.oc.as_ref().unwrap();
1786            let omc = &oc[0];
1787            let orc = &omc.orc.as_ref().unwrap()[0];
1788            let uo = &orc.uo.as_ref().unwrap()[0];
1789
1790            assert!(is_resting_sp_bet(uo));
1791
1792            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
1793            let report = parse_order_status_report(
1794                uo,
1795                instrument_id,
1796                AccountId::from("BETFAIR-001"),
1797                parse_millis_timestamp(ocm.pt),
1798                parse_millis_timestamp(ocm.pt),
1799            )
1800            .unwrap();
1801
1802            // A resting SP bet reports execution-complete with zero size
1803            // fields, but cannot be cancelled and is still open until BSP
1804            // reconciliation
1805            assert_eq!(report.order_status, OrderStatus::Accepted);
1806            assert_eq!(report.time_in_force, TimeInForce::AtTheClose);
1807            assert_eq!(report.order_side, Some(OrderSide::Sell)); // Back → Sell
1808            assert_eq!(report.filled_qty.as_f64(), 0.0);
1809            assert_eq!(report.quantity.as_f64(), 2.0);
1810        } else {
1811            panic!("expected OrderChange");
1812        }
1813    }
1814
1815    fn sp_unmatched_order(size_matched: Decimal, size_lapsed: Decimal) -> UnmatchedOrder {
1816        UnmatchedOrder {
1817            id: "442849719274".to_string(),
1818            p: Decimal::new(10000, 1),
1819            s: Decimal::ZERO,
1820            side: StreamingSide::Back,
1821            status: StreamingOrderStatus::ExecutionComplete,
1822            pt: Some(StreamingPersistenceType::MarketOnClose),
1823            ot: StreamingOrderType::MarketOnClose,
1824            pd: 1789423573000,
1825            bsp: Some(Decimal::new(2, 0)),
1826            rfo: None,
1827            rfs: None,
1828            rc: None,
1829            rac: None,
1830            md: None,
1831            cd: None,
1832            ld: None,
1833            avp: None,
1834            sm: Some(size_matched),
1835            sr: Some(Decimal::ZERO),
1836            sl: Some(size_lapsed),
1837            sc: Some(Decimal::ZERO),
1838            sv: Some(Decimal::ZERO),
1839            lsrc: None,
1840        }
1841    }
1842
1843    #[rstest]
1844    fn test_parse_order_status_report_sp_reconciled() {
1845        // At BSP reconciliation an SP bet is no longer resting: a matched bet
1846        // resolves Filled and a lapsed bet resolves Canceled
1847        let matched = sp_unmatched_order(Decimal::new(2, 0), Decimal::ZERO);
1848        assert!(!is_resting_sp_bet(&matched));
1849
1850        let report = parse_order_status_report(
1851            &matched,
1852            InstrumentId::from("1.262362241-6532924-0.BETFAIR"),
1853            AccountId::from("BETFAIR-001"),
1854            UnixNanos::default(),
1855            UnixNanos::default(),
1856        )
1857        .unwrap();
1858
1859        assert_eq!(report.order_status, OrderStatus::Filled);
1860        assert_eq!(report.filled_qty.as_f64(), 2.0);
1861        assert_eq!(report.quantity.as_f64(), 2.0);
1862
1863        let lapsed = sp_unmatched_order(Decimal::ZERO, Decimal::new(2, 0));
1864        assert!(!is_resting_sp_bet(&lapsed));
1865
1866        let report = parse_order_status_report(
1867            &lapsed,
1868            InstrumentId::from("1.262362241-6532924-0.BETFAIR"),
1869            AccountId::from("BETFAIR-001"),
1870            UnixNanos::default(),
1871            UnixNanos::default(),
1872        )
1873        .unwrap();
1874
1875        assert_eq!(report.order_status, OrderStatus::Canceled);
1876    }
1877
1878    #[rstest]
1879    fn test_parse_order_status_report_multiple_fills() {
1880        let data = load_test_json("stream/ocm_multiple_fills.json");
1881        let msgs: Vec<StreamMessage> = serde_json::from_str(&data).unwrap();
1882
1883        if let StreamMessage::OrderChange(ocm) = &msgs[0] {
1884            let oc = ocm.oc.as_ref().unwrap();
1885            let omc = &oc[0];
1886            let orc = &omc.orc.as_ref().unwrap()[0];
1887            let uo = &orc.uo.as_ref().unwrap()[0];
1888
1889            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
1890            let report = parse_order_status_report(
1891                uo,
1892                instrument_id,
1893                AccountId::from("BETFAIR-001"),
1894                parse_millis_timestamp(ocm.pt),
1895                parse_millis_timestamp(ocm.pt),
1896            )
1897            .unwrap();
1898
1899            // Partially filled: sm=16.19, status=E, has rfo
1900            assert_eq!(report.order_status, OrderStatus::PartiallyFilled);
1901            assert_eq!(report.order_side, Some(OrderSide::Sell)); // Back → Sell
1902            assert_eq!(report.filled_qty.as_f64(), 16.19);
1903            assert!(report.client_order_id.is_some());
1904            assert!(report.avg_px.is_some());
1905        } else {
1906            panic!("expected OrderChange");
1907        }
1908    }
1909
1910    #[rstest]
1911    fn test_parse_runner_book_snapshot_empty_book() {
1912        // A snapshot with no levels should still produce a clear delta
1913        let rc = RunnerChange {
1914            id: 12345,
1915            hc: None,
1916            atb: Some(vec![]),
1917            atl: Some(vec![]),
1918            batb: None,
1919            batl: None,
1920            bdatb: None,
1921            bdatl: None,
1922            spb: None,
1923            spl: None,
1924            spn: None,
1925            spf: None,
1926            trd: None,
1927            ltp: None,
1928            tv: None,
1929        };
1930
1931        let result = parse_runner_book_deltas(
1932            make_instrument_id("1.123", 12345, Decimal::ZERO),
1933            &rc,
1934            true,
1935            1000,
1936            UnixNanos::default(),
1937            UnixNanos::default(),
1938        )
1939        .unwrap()
1940        .expect("should produce snapshot deltas");
1941
1942        // Just the Clear delta with F_LAST
1943        assert_eq!(result.deltas.len(), 1);
1944        assert_eq!(result.deltas[0].action, BookAction::Clear);
1945        assert!(RecordFlag::F_LAST.matches(result.deltas[0].flags));
1946        assert!(RecordFlag::F_SNAPSHOT.matches(result.deltas[0].flags));
1947    }
1948
1949    #[rstest]
1950    fn test_make_fill_report() {
1951        let instrument_id = make_instrument_id("1.180604981", 1209555, Decimal::ZERO);
1952        let fill = make_fill_report(
1953            AccountId::from("BETFAIR-001"),
1954            instrument_id,
1955            VenueOrderId::from("229430281339"),
1956            TradeId::from("229430281339-0"),
1957            OrderSide::Buy,
1958            Quantity::new(10.0, BETFAIR_QUANTITY_PRECISION),
1959            Price::new(1.1, BETFAIR_PRICE_PRECISION),
1960            Currency::GBP(),
1961            None,
1962            UnixNanos::default(),
1963            UnixNanos::default(),
1964        );
1965
1966        assert_eq!(fill.instrument_id, instrument_id);
1967        assert_eq!(fill.order_side, OrderSide::Buy);
1968        assert_eq!(fill.last_qty.as_f64(), 10.0);
1969        assert_eq!(fill.last_px.as_f64(), 1.1);
1970        assert_eq!(fill.commission.as_f64(), 0.0);
1971        assert_eq!(fill.liquidity_side, LiquiditySide::NoLiquiditySide);
1972    }
1973
1974    #[rstest]
1975    fn test_fill_tracker_single_full_fill() {
1976        let data = load_test_json("stream/ocm_FILLED.json");
1977        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
1978
1979        if let StreamMessage::OrderChange(ocm) = msg {
1980            let oc = ocm.oc.as_ref().unwrap();
1981            let omc = &oc[0];
1982            let orc = &omc.orc.as_ref().unwrap()[0];
1983            let uo = &orc.uo.as_ref().unwrap()[0];
1984            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
1985            let ts = parse_millis_timestamp(ocm.pt);
1986
1987            let mut tracker = FillTracker::new();
1988            let fill = tracker
1989                .maybe_fill_report(
1990                    uo,
1991                    uo.s,
1992                    instrument_id,
1993                    AccountId::from("BETFAIR-001"),
1994                    Currency::GBP(),
1995                    ts,
1996                    ts,
1997                )
1998                .expect("should produce fill");
1999
2000            assert_eq!(fill.last_qty.as_f64(), 10.0);
2001            assert_eq!(fill.last_px.as_f64(), 1.1);
2002            assert_eq!(fill.order_side, OrderSide::Buy);
2003            assert!(fill.client_order_id.is_some());
2004        } else {
2005            panic!("expected OrderChange");
2006        }
2007    }
2008
2009    #[rstest]
2010    fn test_fill_tracker_incremental_fills() {
2011        let data = load_test_json("stream/ocm_multiple_fills.json");
2012        let msgs: Vec<StreamMessage> = serde_json::from_str(&data).unwrap();
2013        let instrument_id = make_instrument_id("1.179082386", 50210, Decimal::ZERO);
2014
2015        let mut tracker = FillTracker::new();
2016        let account_id = AccountId::from("BETFAIR-001");
2017        let currency = Currency::GBP();
2018
2019        // First fill: sm=16.19 (from zero)
2020        let uo1 = extract_uo(&msgs[0]);
2021        let ts1 = extract_ts(&msgs[0]);
2022        let fill1 = tracker
2023            .maybe_fill_report(uo1, uo1.s, instrument_id, account_id, currency, ts1, ts1)
2024            .expect("should produce first fill");
2025        assert_eq!(fill1.last_qty.as_f64(), 16.19);
2026        assert_eq!(fill1.last_px.as_f64(), 5.8);
2027
2028        // Second fill: sm=16.96 (delta=0.77)
2029        let uo2 = extract_uo(&msgs[1]);
2030        let ts2 = extract_ts(&msgs[1]);
2031        let fill2 = tracker
2032            .maybe_fill_report(uo2, uo2.s, instrument_id, account_id, currency, ts2, ts2)
2033            .expect("should produce second fill");
2034        assert_eq!(fill2.last_qty.as_f64(), 0.77);
2035        assert_eq!(fill2.last_px.as_f64(), 5.8);
2036
2037        // Third fill: sm=17.73 (delta=0.77)
2038        let uo3 = extract_uo(&msgs[2]);
2039        let ts3 = extract_ts(&msgs[2]);
2040        let fill3 = tracker
2041            .maybe_fill_report(uo3, uo3.s, instrument_id, account_id, currency, ts3, ts3)
2042            .expect("should produce third fill");
2043        assert_eq!(fill3.last_qty.as_f64(), 0.77);
2044    }
2045
2046    #[rstest]
2047    fn test_fill_tracker_different_price() {
2048        let data = load_test_json("stream/ocm_filled_different_price.json");
2049        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2050
2051        if let StreamMessage::OrderChange(ocm) = msg {
2052            let oc = ocm.oc.as_ref().unwrap();
2053            let omc = &oc[0];
2054            let orc = &omc.orc.as_ref().unwrap()[0];
2055            let uo = &orc.uo.as_ref().unwrap()[0];
2056            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
2057            let ts = parse_millis_timestamp(ocm.pt);
2058
2059            let mut tracker = FillTracker::new();
2060            let fill = tracker
2061                .maybe_fill_report(
2062                    uo,
2063                    uo.s,
2064                    instrument_id,
2065                    AccountId::from("BETFAIR-001"),
2066                    Currency::GBP(),
2067                    ts,
2068                    ts,
2069                )
2070                .expect("should produce fill");
2071
2072            // Order placed at 1.3 but average fill price is 1.2
2073            assert_eq!(fill.last_qty.as_f64(), 20.0);
2074            assert_eq!(fill.last_px.as_f64(), 1.2);
2075        } else {
2076            panic!("expected OrderChange");
2077        }
2078    }
2079
2080    #[rstest]
2081    fn test_fill_tracker_cancel_no_fill() {
2082        let data = load_test_json("stream/ocm_CANCEL.json");
2083        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2084
2085        if let StreamMessage::OrderChange(ocm) = msg {
2086            let oc = ocm.oc.as_ref().unwrap();
2087            let omc = &oc[0];
2088            let orc = &omc.orc.as_ref().unwrap()[0];
2089            let uo = &orc.uo.as_ref().unwrap()[0];
2090            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
2091            let ts = parse_millis_timestamp(ocm.pt);
2092
2093            let mut tracker = FillTracker::new();
2094            let result = tracker.maybe_fill_report(
2095                uo,
2096                uo.s,
2097                instrument_id,
2098                AccountId::from("BETFAIR-001"),
2099                Currency::GBP(),
2100                ts,
2101                ts,
2102            );
2103            assert!(result.is_none(), "cancelled order should not produce fill");
2104        } else {
2105            panic!("expected OrderChange");
2106        }
2107    }
2108
2109    #[rstest]
2110    fn test_fill_tracker_lapsed_no_fill() {
2111        let data = load_test_json("stream/ocm_error_fill.json");
2112        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2113
2114        if let StreamMessage::OrderChange(ocm) = msg {
2115            let oc = ocm.oc.as_ref().unwrap();
2116            let omc = &oc[0];
2117            let orc = &omc.orc.as_ref().unwrap()[0];
2118            let uo = &orc.uo.as_ref().unwrap()[0];
2119            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
2120            let ts = parse_millis_timestamp(ocm.pt);
2121
2122            let mut tracker = FillTracker::new();
2123            let result = tracker.maybe_fill_report(
2124                uo,
2125                uo.s,
2126                instrument_id,
2127                AccountId::from("BETFAIR-001"),
2128                Currency::GBP(),
2129                ts,
2130                ts,
2131            );
2132            assert!(result.is_none(), "lapsed order should not produce fill");
2133        } else {
2134            panic!("expected OrderChange");
2135        }
2136    }
2137
2138    #[rstest]
2139    fn test_fill_tracker_duplicate_dedup() {
2140        let data = load_test_json("stream/ocm_FILLED.json");
2141        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2142
2143        if let StreamMessage::OrderChange(ocm) = msg {
2144            let oc = ocm.oc.as_ref().unwrap();
2145            let omc = &oc[0];
2146            let orc = &omc.orc.as_ref().unwrap()[0];
2147            let uo = &orc.uo.as_ref().unwrap()[0];
2148            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
2149            let ts = parse_millis_timestamp(ocm.pt);
2150            let account_id = AccountId::from("BETFAIR-001");
2151            let currency = Currency::GBP();
2152
2153            let mut tracker = FillTracker::new();
2154
2155            // First call produces fill
2156            let fill1 =
2157                tracker.maybe_fill_report(uo, uo.s, instrument_id, account_id, currency, ts, ts);
2158            assert!(fill1.is_some());
2159
2160            // Second call with same data produces nothing (dedup)
2161            let fill2 =
2162                tracker.maybe_fill_report(uo, uo.s, instrument_id, account_id, currency, ts, ts);
2163            assert!(fill2.is_none(), "duplicate fill should be suppressed");
2164        } else {
2165            panic!("expected OrderChange");
2166        }
2167    }
2168
2169    #[rstest]
2170    fn test_fill_tracker_price_back_calculation() {
2171        let data = load_test_json("stream/ocm_multiple_fills.json");
2172        let msgs: Vec<StreamMessage> = serde_json::from_str(&data).unwrap();
2173        let instrument_id = make_instrument_id("1.179082386", 50210, Decimal::ZERO);
2174        let account_id = AccountId::from("BETFAIR-001");
2175        let currency = Currency::GBP();
2176        let mut tracker = FillTracker::new();
2177
2178        // Process first fill at avp=5.8
2179        let uo1 = extract_uo(&msgs[0]);
2180        let ts1 = extract_ts(&msgs[0]);
2181        let fill1 = tracker
2182            .maybe_fill_report(uo1, uo1.s, instrument_id, account_id, currency, ts1, ts1)
2183            .unwrap();
2184        assert_eq!(fill1.last_px.as_f64(), 5.8);
2185
2186        // Second fill also at avp=5.8 (same price, avg unchanged)
2187        let uo2 = extract_uo(&msgs[1]);
2188        let ts2 = extract_ts(&msgs[1]);
2189        let fill2 = tracker
2190            .maybe_fill_report(uo2, uo2.s, instrument_id, account_id, currency, ts2, ts2)
2191            .unwrap();
2192        assert_eq!(fill2.last_px.as_f64(), 5.8);
2193    }
2194
2195    #[rstest]
2196    fn test_has_cancel_quantity() {
2197        let data = load_test_json("stream/ocm_CANCEL.json");
2198        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2199
2200        if let StreamMessage::OrderChange(ocm) = msg {
2201            let uo = &ocm.oc.as_ref().unwrap()[0].orc.as_ref().unwrap()[0]
2202                .uo
2203                .as_ref()
2204                .unwrap()[0];
2205            assert!(has_cancel_quantity(uo));
2206        } else {
2207            panic!("expected OrderChange");
2208        }
2209    }
2210
2211    #[rstest]
2212    fn test_has_cancel_quantity_filled_order() {
2213        let data = load_test_json("stream/ocm_FILLED.json");
2214        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2215
2216        if let StreamMessage::OrderChange(ocm) = msg {
2217            let uo = &ocm.oc.as_ref().unwrap()[0].orc.as_ref().unwrap()[0]
2218                .uo
2219                .as_ref()
2220                .unwrap()[0];
2221            assert!(!has_cancel_quantity(uo));
2222        } else {
2223            panic!("expected OrderChange");
2224        }
2225    }
2226
2227    fn extract_uo(msg: &StreamMessage) -> &UnmatchedOrder {
2228        if let StreamMessage::OrderChange(ocm) = msg {
2229            &ocm.oc.as_ref().unwrap()[0].orc.as_ref().unwrap()[0]
2230                .uo
2231                .as_ref()
2232                .unwrap()[0]
2233        } else {
2234            panic!("expected OrderChange")
2235        }
2236    }
2237
2238    fn extract_ts(msg: &StreamMessage) -> UnixNanos {
2239        if let StreamMessage::OrderChange(ocm) = msg {
2240            parse_millis_timestamp(ocm.pt)
2241        } else {
2242            panic!("expected OrderChange")
2243        }
2244    }
2245
2246    #[rstest]
2247    fn test_parse_race_runner_data_from_fixture() {
2248        let data = load_test_json("stream/rcm_single.json");
2249        let msg = stream_decode(data.as_bytes()).unwrap();
2250
2251        let StreamMessage::RaceChange(rcm) = msg else {
2252            panic!("expected RaceChange");
2253        };
2254
2255        let race = &rcm.rc.as_ref().unwrap()[0];
2256        let rrc = &race.rrc.as_ref().unwrap()[0];
2257        let ts = parse_millis_timestamp(rcm.pt);
2258
2259        let runner = parse_race_runner_data(
2260            race.id.as_deref().unwrap(),
2261            race.mid.as_deref().unwrap(),
2262            rrc,
2263            ts,
2264            ts,
2265        )
2266        .unwrap();
2267
2268        assert_eq!(runner.race_id, "28587288.1650");
2269        assert_eq!(runner.market_id, "1.1234567");
2270        assert_eq!(runner.selection_id, 7390417);
2271        assert!((runner.latitude - 51.4189543).abs() < 1e-6);
2272        assert!((runner.longitude - (-0.4058491)).abs() < 1e-6);
2273        assert!((runner.speed - 17.8).abs() < 1e-6);
2274        assert!((runner.progress - 2051.0).abs() < 1e-6);
2275        assert!((runner.stride_frequency - 2.07).abs() < 1e-6);
2276    }
2277
2278    #[rstest]
2279    fn test_parse_race_progress_from_fixture() {
2280        let data = load_test_json("stream/rcm_single.json");
2281        let msg = stream_decode(data.as_bytes()).unwrap();
2282
2283        let StreamMessage::RaceChange(rcm) = msg else {
2284            panic!("expected RaceChange");
2285        };
2286
2287        let race = &rcm.rc.as_ref().unwrap()[0];
2288        let rpc = race.rpc.as_ref().unwrap();
2289        let ts = parse_millis_timestamp(rcm.pt);
2290
2291        let progress = parse_race_progress(
2292            race.id.as_deref().unwrap(),
2293            race.mid.as_deref().unwrap(),
2294            rpc,
2295            ts,
2296            ts,
2297        );
2298
2299        assert_eq!(progress.race_id, "28587288.1650");
2300        assert_eq!(progress.market_id, "1.1234567");
2301        assert_eq!(progress.gate_name, "1f");
2302        assert!((progress.sectional_time - 10.6).abs() < 1e-6);
2303        assert!((progress.running_time - 46.7).abs() < 1e-6);
2304        assert!((progress.speed - 17.8).abs() < 1e-6);
2305        assert!((progress.progress - 87.5).abs() < 1e-6);
2306
2307        let order: Vec<i64> = serde_json::from_str(&progress.order).unwrap();
2308        assert_eq!(order, vec![7390417, 5600338, 11527189, 6395118, 8706072]);
2309
2310        let jumps: Vec<serde_json::Value> = serde_json::from_str(&progress.jumps).unwrap();
2311        assert_eq!(jumps.len(), 2);
2312        assert_eq!(jumps[0]["J"], 2);
2313    }
2314
2315    #[rstest]
2316    fn test_parse_cricket_match_from_fixture() {
2317        let data = load_test_json("stream/ccm_single.json");
2318        let msg = stream_decode(data.as_bytes()).unwrap();
2319
2320        let StreamMessage::CricketChange(ccm) = msg else {
2321            panic!("expected CricketChange");
2322        };
2323
2324        let change = &ccm.cc.as_ref().unwrap()[0];
2325        let ts = parse_millis_timestamp(ccm.pt);
2326        let cricket = parse_cricket_match(change, ts, ts).unwrap();
2327
2328        assert_eq!(cricket.event_id, "35741575");
2329        assert_eq!(cricket.market_id, "1.259334639");
2330        let stats: serde_json::Value = serde_json::from_str(&cricket.match_stats).unwrap();
2331        assert_eq!(stats["inningsStats"][0]["inningsRuns"], 100);
2332    }
2333
2334    #[rstest]
2335    #[case(None, Some("1.259334639".to_string()))]
2336    #[case(Some("35741575".to_string()), None)]
2337    fn test_parse_cricket_match_missing_required_id_returns_none(
2338        #[case] event_id: Option<String>,
2339        #[case] market_id: Option<String>,
2340    ) {
2341        let change = CricketChange {
2342            event_id,
2343            market_id,
2344            fixture_info: None,
2345            home_team: None,
2346            away_team: None,
2347            match_stats: Some(serde_json::json!({"inningsStats": []})),
2348            incident_list_wrapper: None,
2349        };
2350        let ts = UnixNanos::from(1_000_000_000u64);
2351        let result = parse_cricket_match(&change, ts, ts);
2352
2353        assert!(result.is_none());
2354    }
2355
2356    #[rstest]
2357    fn test_parse_race_runner_data_multi_runner() {
2358        let data = load_test_json("stream/rcm_multi_runner.json");
2359        let msg = stream_decode(data.as_bytes()).unwrap();
2360
2361        let StreamMessage::RaceChange(rcm) = msg else {
2362            panic!("expected RaceChange");
2363        };
2364
2365        let race = &rcm.rc.as_ref().unwrap()[0];
2366        let ts = parse_millis_timestamp(rcm.pt);
2367        let race_id = race.id.as_deref().unwrap();
2368        let market_id = race.mid.as_deref().unwrap();
2369
2370        let runners: Vec<_> = race
2371            .rrc
2372            .as_ref()
2373            .unwrap()
2374            .iter()
2375            .filter_map(|rrc| parse_race_runner_data(race_id, market_id, rrc, ts, ts))
2376            .collect();
2377
2378        assert_eq!(runners.len(), 5);
2379        assert_eq!(runners[0].selection_id, 35467839);
2380        assert_eq!(runners[4].selection_id, 41694785);
2381        assert!((runners[0].speed - 16.33).abs() < 1e-6);
2382        assert!((runners[4].speed - 17.11).abs() < 1e-6);
2383    }
2384
2385    #[rstest]
2386    fn test_parse_race_runner_data_missing_id_returns_none() {
2387        let rrc = RaceRunnerChange {
2388            ft: Some(1000),
2389            id: None,
2390            lat: Some(51.0),
2391            lng: Some(-0.4),
2392            spd: Some(15.0),
2393            prg: Some(500.0),
2394            sfq: Some(2.0),
2395        };
2396        let ts = UnixNanos::from(1_000_000_000u64);
2397        let result = parse_race_runner_data("race1", "market1", &rrc, ts, ts);
2398        assert!(result.is_none());
2399    }
2400
2401    #[rstest]
2402    fn test_parse_race_runner_data_absent_fields_are_nan() {
2403        let rrc = RaceRunnerChange {
2404            ft: None,
2405            id: Some(12345),
2406            lat: None,
2407            lng: None,
2408            spd: None,
2409            prg: None,
2410            sfq: None,
2411        };
2412        let ts = UnixNanos::from(1_000_000_000u64);
2413        let runner = parse_race_runner_data("race1", "market1", &rrc, ts, ts).unwrap();
2414        assert!(runner.latitude.is_nan());
2415        assert!(runner.longitude.is_nan());
2416        assert!(runner.speed.is_nan());
2417        assert!(runner.progress.is_nan());
2418        assert!(runner.stride_frequency.is_nan());
2419    }
2420
2421    #[rstest]
2422    fn test_parse_race_progress_absent_fields() {
2423        let rpc = RaceProgressChange {
2424            ft: None,
2425            g: None,
2426            st: None,
2427            rt: None,
2428            spd: None,
2429            prg: None,
2430            ord: None,
2431            jumps: None,
2432        };
2433        let ts = UnixNanos::from(1_000_000_000u64);
2434        let progress = parse_race_progress("race1", "market1", &rpc, ts, ts);
2435        assert_eq!(progress.gate_name, "");
2436        assert!(progress.sectional_time.is_nan());
2437        assert!(progress.running_time.is_nan());
2438        assert_eq!(progress.order, "");
2439        assert_eq!(progress.jumps, "");
2440    }
2441
2442    fn runner_change_with_ticker(
2443        id: u64,
2444        ltp: Option<Decimal>,
2445        tv: Option<Decimal>,
2446        spn: Option<Decimal>,
2447        spf: Option<Decimal>,
2448    ) -> RunnerChange {
2449        RunnerChange {
2450            id,
2451            hc: None,
2452            atb: None,
2453            atl: None,
2454            batb: None,
2455            batl: None,
2456            bdatb: None,
2457            bdatl: None,
2458            spb: None,
2459            spl: None,
2460            spn,
2461            spf,
2462            trd: None,
2463            ltp,
2464            tv,
2465        }
2466    }
2467
2468    #[rstest]
2469    fn test_parse_betfair_ticker_all_fields() {
2470        let rc = runner_change_with_ticker(
2471            9249757,
2472            Some(Decimal::new(55, 1)),
2473            Some(Decimal::new(189032, 2)),
2474            Some(Decimal::new(568, 2)),
2475            Some(Decimal::new(573, 2)),
2476        );
2477        let ts = UnixNanos::from(1_000_000_000u64);
2478        let instrument_id = make_instrument_id("1.185781465", 9249757, Decimal::ZERO);
2479
2480        let ticker = parse_betfair_ticker(instrument_id, &rc, ts, ts).unwrap();
2481
2482        assert_eq!(ticker.instrument_id, instrument_id);
2483        assert_decimal_option_eq(ticker.last_traded_price, Decimal::new(55, 1));
2484        assert_decimal_option_eq(ticker.traded_volume, Decimal::new(189032, 2));
2485        assert_decimal_option_eq(ticker.starting_price_near, Decimal::new(568, 2));
2486        assert_decimal_option_eq(ticker.starting_price_far, Decimal::new(573, 2));
2487    }
2488
2489    #[rstest]
2490    fn test_parse_betfair_ticker_partial_fields() {
2491        let rc = runner_change_with_ticker(
2492            9249757,
2493            Some(Decimal::new(55, 1)),
2494            Some(Decimal::new(189032, 2)),
2495            None,
2496            None,
2497        );
2498        let ts = UnixNanos::from(1_000_000_000u64);
2499        let instrument_id = make_instrument_id("1.185781465", 9249757, Decimal::ZERO);
2500
2501        let ticker = parse_betfair_ticker(instrument_id, &rc, ts, ts).unwrap();
2502
2503        assert_decimal_option_eq(ticker.last_traded_price, Decimal::new(55, 1));
2504        assert_decimal_option_eq(ticker.traded_volume, Decimal::new(189032, 2));
2505        assert!(ticker.starting_price_near.is_none());
2506        assert!(ticker.starting_price_far.is_none());
2507    }
2508
2509    #[rstest]
2510    fn test_parse_betfair_ticker_no_fields_returns_none() {
2511        let rc = runner_change_with_ticker(9249757, None, None, None, None);
2512        let ts = UnixNanos::from(1_000_000_000u64);
2513        let instrument_id = make_instrument_id("1.185781465", 9249757, Decimal::ZERO);
2514
2515        assert!(parse_betfair_ticker(instrument_id, &rc, ts, ts).is_none());
2516    }
2517
2518    #[rstest]
2519    fn test_parse_betfair_ticker_only_tv() {
2520        let rc =
2521            runner_change_with_ticker(40273293, None, Some(Decimal::new(320115, 2)), None, None);
2522        let ts = UnixNanos::from(1_000_000_000u64);
2523        let instrument_id = make_instrument_id("1.185781465", 40273293, Decimal::ZERO);
2524
2525        let ticker = parse_betfair_ticker(instrument_id, &rc, ts, ts).unwrap();
2526
2527        assert!(ticker.last_traded_price.is_none());
2528        assert_decimal_option_eq(ticker.traded_volume, Decimal::new(320115, 2));
2529        assert!(ticker.starting_price_near.is_none());
2530        assert!(ticker.starting_price_far.is_none());
2531    }
2532
2533    #[rstest]
2534    fn test_parse_betfair_ticker_from_fixture() {
2535        let data = load_test_json("stream/mcm_BSP_settled.json");
2536        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2537
2538        if let StreamMessage::MarketChange(mcm) = msg {
2539            let mc = mcm.mc.as_ref().unwrap();
2540            let change = &mc[0];
2541            let rc_list = change.rc.as_ref().unwrap();
2542
2543            // Runner 9249757 has ltp=5.5, tv=1890.32, spn=5.68, spf=5.73
2544            let rc = rc_list.iter().find(|r| r.id == 9249757).unwrap();
2545            let instrument_id = make_instrument_id(&change.id, rc.id, Decimal::ZERO);
2546            let ts = parse_millis_timestamp(mcm.pt);
2547
2548            let ticker = parse_betfair_ticker(instrument_id, rc, ts, ts).unwrap();
2549
2550            assert_decimal_option_eq(ticker.last_traded_price, Decimal::new(55, 1));
2551            assert_decimal_option_eq(ticker.traded_volume, Decimal::new(189032, 2));
2552            assert_decimal_option_eq(ticker.starting_price_near, Decimal::new(568, 2));
2553            assert_decimal_option_eq(ticker.starting_price_far, Decimal::new(573, 2));
2554
2555            // Runner 40273293 has ltp=2.1, tv=3201.15 but no spn/spf
2556            let rc2 = rc_list.iter().find(|r| r.id == 40273293).unwrap();
2557            let instrument_id2 = make_instrument_id(&change.id, rc2.id, Decimal::ZERO);
2558            let ticker2 = parse_betfair_ticker(instrument_id2, rc2, ts, ts).unwrap();
2559
2560            assert_decimal_option_eq(ticker2.last_traded_price, Decimal::new(21, 1));
2561            assert_decimal_option_eq(ticker2.traded_volume, Decimal::new(320115, 2));
2562            assert!(ticker2.starting_price_near.is_none());
2563            assert!(ticker2.starting_price_far.is_none());
2564
2565            // Runner 23678734 has only tv=0, no ltp
2566            let rc3 = rc_list.iter().find(|r| r.id == 23678734).unwrap();
2567            let instrument_id3 = make_instrument_id(&change.id, rc3.id, Decimal::ZERO);
2568            let ticker3 = parse_betfair_ticker(instrument_id3, rc3, ts, ts).unwrap();
2569
2570            assert!(ticker3.last_traded_price.is_none());
2571            assert_decimal_option_eq(ticker3.traded_volume, Decimal::ZERO);
2572        } else {
2573            panic!("Expected MarketChange");
2574        }
2575    }
2576
2577    #[rstest]
2578    fn test_parse_betfair_starting_prices_from_fixture() {
2579        let data = load_test_json("stream/mcm_BSP_settled.json");
2580        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2581
2582        if let StreamMessage::MarketChange(mcm) = msg {
2583            let mc = mcm.mc.as_ref().unwrap();
2584            let def = mc[0].market_definition.as_ref().unwrap();
2585            let ts = parse_millis_timestamp(mcm.pt);
2586
2587            let prices = parse_betfair_starting_prices(&mc[0].id, def, ts, ts);
2588
2589            // 3 runners have bsp values, 1 (REMOVED) does not
2590            assert_eq!(prices.len(), 3);
2591
2592            let bsp_map: std::collections::HashMap<String, Decimal> = prices
2593                .iter()
2594                .map(|p| (p.instrument_id.to_string(), p.bsp))
2595                .collect();
2596
2597            let id_winner = make_instrument_id("1.185781465", 9249757, Decimal::ZERO).to_string();
2598            let id_placed = make_instrument_id("1.185781465", 40273293, Decimal::ZERO).to_string();
2599            let id_loser = make_instrument_id("1.185781465", 11120000, Decimal::ZERO).to_string();
2600
2601            assert_eq!(bsp_map[&id_winner], Decimal::new(573, 2));
2602            assert_eq!(bsp_map[&id_placed], Decimal::new(214, 2));
2603            assert_eq!(bsp_map[&id_loser], Decimal::new(2856, 2));
2604        } else {
2605            panic!("Expected MarketChange");
2606        }
2607    }
2608
2609    #[rstest]
2610    fn test_parse_betfair_starting_prices_no_runners() {
2611        let def = MarketDefinition {
2612            runners: None,
2613            bet_delay: None,
2614            betting_type: None,
2615            bsp_market: None,
2616            bsp_reconciled: None,
2617            competition_id: None,
2618            competition_name: None,
2619            complete: None,
2620            country_code: None,
2621            cross_matching: None,
2622            discount_allowed: None,
2623            each_way_divisor: None,
2624            event_id: None,
2625            event_name: None,
2626            event_type_id: None,
2627            event_type_name: None,
2628            in_play: None,
2629            line_interval: None,
2630            line_max_unit: None,
2631            line_min_unit: None,
2632            market_base_rate: None,
2633            market_id: None,
2634            market_name: None,
2635            market_time: None,
2636            market_type: None,
2637            number_of_active_runners: None,
2638            number_of_winners: None,
2639            open_date: None,
2640            persistence_enabled: None,
2641            price_ladder_definition: None,
2642            race_type: None,
2643            regulators: None,
2644            runners_voidable: None,
2645            settled_time: None,
2646            status: None,
2647            suspend_time: None,
2648            timezone: None,
2649            turn_in_play_enabled: None,
2650            venue: None,
2651            version: None,
2652        };
2653        let ts = UnixNanos::from(1_000_000_000u64);
2654
2655        let prices = parse_betfair_starting_prices("1.12345", &def, ts, ts);
2656
2657        assert!(prices.is_empty());
2658    }
2659
2660    #[rstest]
2661    fn test_parse_bsp_book_deltas_from_fixture() {
2662        let data = load_test_json("stream/mcm_BSP.json");
2663        let messages: Vec<StreamMessage> = serde_json::from_str(&data).unwrap();
2664
2665        // Find the MCM with runner changes containing spb/spl data
2666        let mcm = messages
2667            .iter()
2668            .find_map(|m| match m {
2669                StreamMessage::MarketChange(mcm) => {
2670                    let mc = mcm.mc.as_ref()?;
2671                    let has_spb = mc.iter().any(|c| {
2672                        c.rc.as_ref()
2673                            .is_some_and(|rcs| rcs.iter().any(|r| r.spb.is_some()))
2674                    });
2675
2676                    if has_spb { Some(mcm) } else { None }
2677                }
2678                _ => None,
2679            })
2680            .expect("fixture should contain MCM with spb data");
2681
2682        let mc = mcm.mc.as_ref().unwrap();
2683        let change = &mc[0];
2684        let rc_list = change.rc.as_ref().unwrap();
2685
2686        // Runner 9249757 has spb and spl arrays
2687        let rc = rc_list.iter().find(|r| r.id == 9249757).unwrap();
2688        let instrument_id = make_instrument_id(&change.id, rc.id, Decimal::ZERO);
2689        let ts = parse_millis_timestamp(mcm.pt);
2690
2691        let deltas = parse_bsp_book_deltas(instrument_id, rc, ts, ts);
2692
2693        let spb_count = rc.spb.as_ref().unwrap().len();
2694        let spl_count = rc.spl.as_ref().unwrap().len();
2695        assert_eq!(deltas.len(), spb_count + spl_count);
2696
2697        // SPB entries are Sell side
2698        assert_eq!(deltas[0].side, OrderSide::Sell);
2699        assert_eq!(deltas[0].price, Decimal::new(1000, 0));
2700        assert_eq!(deltas[0].size, Decimal::new(3338, 2));
2701        assert_eq!(deltas[0].action, BookAction::Update);
2702
2703        // SPL entries are Buy side
2704        let spl_start = spb_count;
2705        assert_eq!(deltas[spl_start].side, OrderSide::Buy);
2706        assert_eq!(deltas[spl_start].price, Decimal::new(7, 0));
2707        assert_eq!(deltas[spl_start].size, Decimal::new(10, 0));
2708    }
2709
2710    #[rstest]
2711    fn test_parse_bsp_book_deltas_zero_volume_is_delete() {
2712        let rc = RunnerChange {
2713            id: 12345,
2714            hc: None,
2715            atb: None,
2716            atl: None,
2717            batb: None,
2718            batl: None,
2719            bdatb: None,
2720            bdatl: None,
2721            spb: Some(vec![PV {
2722                price: Decimal::new(50, 1),
2723                volume: Decimal::ZERO,
2724            }]),
2725            spl: None,
2726            spn: None,
2727            spf: None,
2728            trd: None,
2729            ltp: None,
2730            tv: None,
2731        };
2732        let ts = UnixNanos::from(1_000_000_000u64);
2733        let instrument_id = make_instrument_id("1.12345", 12345, Decimal::ZERO);
2734
2735        let deltas = parse_bsp_book_deltas(instrument_id, &rc, ts, ts);
2736
2737        assert_eq!(deltas.len(), 1);
2738        assert_eq!(deltas[0].action, BookAction::Delete);
2739        assert_eq!(deltas[0].price, Decimal::new(5, 0));
2740        assert_eq!(deltas[0].size, Decimal::ZERO);
2741    }
2742
2743    #[rstest]
2744    fn test_parse_bsp_book_deltas_no_spb_spl_returns_empty() {
2745        let rc = runner_change_with_ticker(12345, None, None, None, None);
2746        let ts = UnixNanos::from(1_000_000_000u64);
2747        let instrument_id = make_instrument_id("1.12345", 12345, Decimal::ZERO);
2748
2749        let deltas = parse_bsp_book_deltas(instrument_id, &rc, ts, ts);
2750
2751        assert!(deltas.is_empty());
2752    }
2753
2754    #[rstest]
2755    fn test_parse_instrument_closes_from_fixture() {
2756        let data = load_test_json("stream/mcm_BSP_settled.json");
2757        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
2758
2759        if let StreamMessage::MarketChange(mcm) = msg {
2760            let mc = mcm.mc.as_ref().unwrap();
2761            let def = mc[0].market_definition.as_ref().unwrap();
2762            let ts = parse_millis_timestamp(mcm.pt);
2763
2764            let closes = parse_instrument_closes(&mc[0].id, def, ts, ts);
2765
2766            // 4 runners: WINNER, PLACED, LOSER, REMOVED - all produce close events
2767            assert_eq!(closes.len(), 4);
2768
2769            let close_map: std::collections::HashMap<String, Price> = closes
2770                .iter()
2771                .map(|c| (c.instrument_id.to_string(), c.close_price))
2772                .collect();
2773
2774            let id_winner = make_instrument_id("1.185781465", 9249757, Decimal::ZERO).to_string();
2775            let id_placed = make_instrument_id("1.185781465", 40273293, Decimal::ZERO).to_string();
2776            let id_loser = make_instrument_id("1.185781465", 11120000, Decimal::ZERO).to_string();
2777            let id_removed = make_instrument_id("1.185781465", 37433527, Decimal::ZERO).to_string();
2778
2779            assert_eq!(close_map[&id_winner], Price::from("1.00"));
2780            assert_eq!(close_map[&id_placed], Price::from("1.00"));
2781            assert_eq!(close_map[&id_loser], Price::from("0.00"));
2782            assert_eq!(close_map[&id_removed], Price::from("0.00"));
2783        } else {
2784            panic!("Expected MarketChange");
2785        }
2786    }
2787
2788    #[rstest]
2789    fn test_parse_instrument_closes_active_runners_excluded() {
2790        let data = load_test_json("stream/mcm_BSP.json");
2791        let messages: Vec<StreamMessage> = serde_json::from_str(&data).unwrap();
2792
2793        // Find the MCM with a market definition containing ACTIVE runners
2794        let mcm = messages
2795            .iter()
2796            .find_map(|m| match m {
2797                StreamMessage::MarketChange(mcm) => {
2798                    let mc = mcm.mc.as_ref()?;
2799                    mc.iter()
2800                        .find(|c| c.market_definition.is_some())
2801                        .map(|_| mcm)
2802                }
2803                _ => None,
2804            })
2805            .expect("fixture should contain MCM with market definition");
2806
2807        let mc = mcm.mc.as_ref().unwrap();
2808        let change = mc.iter().find(|c| c.market_definition.is_some()).unwrap();
2809        let def = change.market_definition.as_ref().unwrap();
2810        let ts = parse_millis_timestamp(mcm.pt);
2811
2812        let closes = parse_instrument_closes(&change.id, def, ts, ts);
2813
2814        assert!(
2815            closes.is_empty(),
2816            "Active runners should not produce close events, found {}",
2817            closes.len()
2818        );
2819    }
2820
2821    fn make_test_uo(
2822        bet_id: &str,
2823        size: Decimal,
2824        sm: Option<Decimal>,
2825        avp: Option<Decimal>,
2826    ) -> UnmatchedOrder {
2827        UnmatchedOrder {
2828            id: bet_id.to_string(),
2829            p: Decimal::new(25, 1),
2830            s: size,
2831            side: StreamingSide::Back,
2832            status: StreamingOrderStatus::Executable,
2833            pt: Some(StreamingPersistenceType::Lapse),
2834            ot: StreamingOrderType::Limit,
2835            pd: 1616568581000,
2836            bsp: None,
2837            rfo: None,
2838            rfs: None,
2839            rc: None,
2840            rac: None,
2841            md: None,
2842            cd: None,
2843            ld: None,
2844            avp,
2845            sm,
2846            sr: None,
2847            sl: None,
2848            sc: None,
2849            sv: None,
2850            lsrc: None,
2851        }
2852    }
2853
2854    #[rstest]
2855    fn test_make_trade_id_normalizes_float_tail_sm() {
2856        let uo = make_test_uo(
2857            "430069890490",
2858            Decimal::new(4287, 2),
2859            Some(Decimal::new(4_287_000_000_000_001, 14)),
2860            Some(Decimal::new(25, 1)),
2861        );
2862
2863        let trade_id = make_trade_id(&uo);
2864
2865        assert_eq!(trade_id.as_str(), "430069890490-42.87");
2866    }
2867
2868    #[rstest]
2869    fn test_fill_tracker_sync_order_prevents_duplicate_fill() {
2870        let mut tracker = FillTracker::new();
2871
2872        // Sync existing fill state: 10 matched at 2.5
2873        tracker.sync_order("123456", Decimal::new(10, 0), Decimal::new(25, 1));
2874
2875        let uo = make_test_uo(
2876            "123456",
2877            Decimal::new(20, 0),
2878            Some(Decimal::new(10, 0)),
2879            Some(Decimal::new(25, 1)),
2880        );
2881
2882        let instrument_id = InstrumentId::from("1.234567-123456-0.0.BETFAIR");
2883        let account_id = AccountId::from("BETFAIR-001");
2884        let currency = Currency::from("GBP");
2885        let ts = UnixNanos::default();
2886
2887        // Same sm=10 as synced, should not emit a fill
2888        let result =
2889            tracker.maybe_fill_report(&uo, uo.s, instrument_id, account_id, currency, ts, ts);
2890        assert!(
2891            result.is_none(),
2892            "should not emit fill for already-synced qty"
2893        );
2894    }
2895
2896    #[rstest]
2897    fn test_fill_tracker_sync_order_is_monotonic() {
2898        let mut tracker = FillTracker::new();
2899
2900        // Tracker already has 10 from prior stream activity
2901        tracker.sync_order("123456", Decimal::new(10, 0), Decimal::new(25, 1));
2902
2903        // Cache lags (engine hasn't processed the last fill yet) and reports 5
2904        tracker.sync_order("123456", Decimal::new(5, 0), Decimal::new(20, 1));
2905
2906        // Next OCM with sm=15 must emit incremental fill of 5, not 10
2907        let uo = make_test_uo(
2908            "123456",
2909            Decimal::new(20, 0),
2910            Some(Decimal::new(15, 0)),
2911            Some(Decimal::new(25, 1)),
2912        );
2913        let result = tracker.maybe_fill_report(
2914            &uo,
2915            uo.s,
2916            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
2917            AccountId::from("BETFAIR-001"),
2918            Currency::from("GBP"),
2919            UnixNanos::default(),
2920            UnixNanos::default(),
2921        );
2922        let fill = result.expect("should emit incremental fill");
2923        assert_eq!(fill.last_qty, Quantity::from("5.00"));
2924    }
2925
2926    #[rstest]
2927    fn test_fill_tracker_sync_order_allows_incremental_fill() {
2928        let mut tracker = FillTracker::new();
2929
2930        // Sync: 10 matched at 2.5
2931        tracker.sync_order("123456", Decimal::new(10, 0), Decimal::new(25, 1));
2932
2933        let uo = make_test_uo(
2934            "123456",
2935            Decimal::new(20, 0),
2936            Some(Decimal::new(15, 0)),
2937            Some(Decimal::new(26, 1)),
2938        );
2939
2940        let instrument_id = InstrumentId::from("1.234567-123456-0.0.BETFAIR");
2941        let account_id = AccountId::from("BETFAIR-001");
2942        let currency = Currency::from("GBP");
2943        let ts = UnixNanos::default();
2944
2945        // sm=15 vs synced 10, should emit incremental fill of 5
2946        let result =
2947            tracker.maybe_fill_report(&uo, uo.s, instrument_id, account_id, currency, ts, ts);
2948        assert!(result.is_some(), "should emit fill for new matched qty");
2949        let fill = result.unwrap();
2950        assert_eq!(fill.last_qty, Quantity::from("5.00"));
2951    }
2952
2953    #[rstest]
2954    fn test_fill_tracker_seed_published_trade_ids_blocks_replay() {
2955        let mut tracker = FillTracker::new();
2956
2957        let uo_for_id = make_test_uo(
2958            "123456",
2959            Decimal::new(20, 0),
2960            Some(Decimal::new(10, 0)),
2961            Some(Decimal::new(25, 1)),
2962        );
2963        let trade_id = make_trade_id(&uo_for_id);
2964
2965        // Seeded trade-id must block even with an empty FillTracker
2966        // (cumulative-size gate would otherwise let this fill through)
2967        tracker.seed_published_trade_ids([trade_id.to_string()]);
2968
2969        let result = tracker.maybe_fill_report(
2970            &uo_for_id,
2971            uo_for_id.s,
2972            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
2973            AccountId::from("BETFAIR-001"),
2974            Currency::from("GBP"),
2975            UnixNanos::default(),
2976            UnixNanos::default(),
2977        );
2978
2979        assert!(
2980            result.is_none(),
2981            "seeded trade-id should suppress duplicate fill",
2982        );
2983
2984        let next_update = make_test_uo(
2985            "123456",
2986            Decimal::new(20, 0),
2987            Some(Decimal::new(15, 0)),
2988            Some(Decimal::new(25, 1)),
2989        );
2990        let next_fill = tracker
2991            .maybe_fill_report(
2992                &next_update,
2993                next_update.s,
2994                InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
2995                AccountId::from("BETFAIR-001"),
2996                Currency::from("GBP"),
2997                UnixNanos::default(),
2998                UnixNanos::default(),
2999            )
3000            .expect("later cumulative state should emit only the new delta");
3001        assert_eq!(next_fill.last_qty, Quantity::from("5.00"));
3002    }
3003
3004    #[rstest]
3005    fn test_fill_tracker_accepts_sm_float_tail_at_order_quantity() {
3006        let mut tracker = FillTracker::new();
3007        let uo = make_test_uo(
3008            "430069890490",
3009            Decimal::new(4287, 2),
3010            Some(Decimal::new(4_287_000_000_000_001, 14)),
3011            Some(Decimal::new(25, 1)),
3012        );
3013
3014        let result = tracker.maybe_fill_report(
3015            &uo,
3016            uo.s,
3017            InstrumentId::from("1.234567-430069890490-0.0.BETFAIR"),
3018            AccountId::from("BETFAIR-001"),
3019            Currency::from("GBP"),
3020            UnixNanos::default(),
3021            UnixNanos::default(),
3022        );
3023
3024        let fill = result.expect("float tail at order quantity should emit fill");
3025        assert_eq!(fill.last_qty, Quantity::from("42.87"));
3026        assert_eq!(fill.trade_id.as_str(), "430069890490-42.87");
3027
3028        let duplicate = tracker.maybe_fill_report(
3029            &uo,
3030            uo.s,
3031            InstrumentId::from("1.234567-430069890490-0.0.BETFAIR"),
3032            AccountId::from("BETFAIR-001"),
3033            Currency::from("GBP"),
3034            UnixNanos::default(),
3035            UnixNanos::default(),
3036        );
3037        assert!(
3038            duplicate.is_none(),
3039            "same normalized sm should not emit a duplicate fill"
3040        );
3041    }
3042
3043    #[rstest]
3044    fn test_fill_tracker_normalized_duplicate_updates_average_price_anchor() {
3045        let mut tracker = FillTracker::new();
3046
3047        let first = tracker.advance_cumulative_fill(
3048            "123456",
3049            Decimal::new(10, 0),
3050            Some(Decimal::new(20, 1)),
3051            Decimal::new(20, 1),
3052        );
3053        let normalized_duplicate = tracker.advance_cumulative_fill(
3054            "123456",
3055            Decimal::new(10001, 3),
3056            Some(Decimal::new(30, 1)),
3057            Decimal::new(30, 1),
3058        );
3059        let next = tracker
3060            .advance_cumulative_fill(
3061                "123456",
3062                Decimal::new(11, 0),
3063                Some(Decimal::new(30, 1)),
3064                Decimal::new(30, 1),
3065            )
3066            .expect("new cumulative quantity should emit a fill");
3067
3068        assert!(first.is_some());
3069        assert!(normalized_duplicate.is_none());
3070        assert_eq!(next.1, Quantity::from("1.00"));
3071        assert_eq!(next.2, Price::from("3.00"));
3072    }
3073
3074    #[rstest]
3075    fn test_fill_tracker_overfill_rejected() {
3076        let mut tracker = FillTracker::new();
3077
3078        // sm=30 exceeds order size s=20
3079        let uo = make_test_uo(
3080            "999001",
3081            Decimal::new(20, 0),
3082            Some(Decimal::new(30, 0)),
3083            Some(Decimal::new(25, 1)),
3084        );
3085
3086        let instrument_id = InstrumentId::from("1.234567-999001-0.0.BETFAIR");
3087        let account_id = AccountId::from("BETFAIR-001");
3088        let currency = Currency::from("GBP");
3089        let ts = UnixNanos::default();
3090
3091        let result =
3092            tracker.maybe_fill_report(&uo, uo.s, instrument_id, account_id, currency, ts, ts);
3093        assert!(
3094            result.is_none(),
3095            "overfill (sm > order_qty) should be rejected"
3096        );
3097    }
3098
3099    #[rstest]
3100    fn test_fill_tracker_zero_sm_returns_none() {
3101        let mut tracker = FillTracker::new();
3102
3103        let uo = make_test_uo("999002", Decimal::new(10, 0), Some(Decimal::ZERO), None);
3104
3105        let instrument_id = InstrumentId::from("1.234567-999002-0.0.BETFAIR");
3106        let result = tracker.maybe_fill_report(
3107            &uo,
3108            uo.s,
3109            instrument_id,
3110            AccountId::from("BETFAIR-001"),
3111            Currency::from("GBP"),
3112            UnixNanos::default(),
3113            UnixNanos::default(),
3114        );
3115        assert!(result.is_none(), "zero sm should not produce a fill");
3116    }
3117
3118    #[rstest]
3119    fn test_fill_tracker_no_avp_uses_order_price() {
3120        let data = load_test_json("stream/ocm_FILLED_no_avp.json");
3121        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
3122
3123        if let StreamMessage::OrderChange(ocm) = msg {
3124            let oc = ocm.oc.as_ref().unwrap();
3125            let omc = &oc[0];
3126            let orc = &omc.orc.as_ref().unwrap()[0];
3127            let uo = &orc.uo.as_ref().unwrap()[0];
3128            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
3129            let ts = parse_millis_timestamp(ocm.pt);
3130
3131            let mut tracker = FillTracker::new();
3132            let fill = tracker
3133                .maybe_fill_report(
3134                    uo,
3135                    uo.s,
3136                    instrument_id,
3137                    AccountId::from("BETFAIR-001"),
3138                    Currency::GBP(),
3139                    ts,
3140                    ts,
3141                )
3142                .expect("should produce fill even without avp");
3143
3144            // No avp field, so fill price falls back to order price (p=3.5)
3145            assert_eq!(fill.last_qty.as_f64(), 25.0);
3146            assert_eq!(fill.last_px.as_f64(), 3.5);
3147        } else {
3148            panic!("expected OrderChange");
3149        }
3150    }
3151
3152    #[rstest]
3153    fn test_fill_tracker_weighted_avg_back_calculation() {
3154        let mut tracker = FillTracker::new();
3155        let instrument_id = InstrumentId::from("1.234567-999003-0.0.BETFAIR");
3156        let account_id = AccountId::from("BETFAIR-001");
3157        let currency = Currency::from("GBP");
3158        let ts = UnixNanos::default();
3159
3160        // First fill: 10 @ avp=2.0
3161        let uo1 = make_test_uo(
3162            "999003",
3163            Decimal::new(30, 0),
3164            Some(Decimal::new(10, 0)),
3165            Some(Decimal::new(20, 1)),
3166        );
3167        let fill1 = tracker
3168            .maybe_fill_report(&uo1, uo1.s, instrument_id, account_id, currency, ts, ts)
3169            .expect("first fill");
3170        assert_eq!(fill1.last_px.as_f64(), 2.0);
3171        assert_eq!(fill1.last_qty.as_f64(), 10.0);
3172
3173        // Second fill: sm=20, avp=2.5
3174        // Back-calc: (2.5*20 - 2.0*10) / 10 = (50-20)/10 = 3.0
3175        let uo2 = make_test_uo(
3176            "999003",
3177            Decimal::new(30, 0),
3178            Some(Decimal::new(20, 0)),
3179            Some(Decimal::new(25, 1)),
3180        );
3181        let fill2 = tracker
3182            .maybe_fill_report(&uo2, uo2.s, instrument_id, account_id, currency, ts, ts)
3183            .expect("second fill");
3184        assert_eq!(fill2.last_qty.as_f64(), 10.0);
3185        assert_eq!(fill2.last_px.as_f64(), 3.0);
3186    }
3187
3188    #[rstest]
3189    fn test_fill_tracker_negative_fill_price_falls_back_to_avp() {
3190        let mut tracker = FillTracker::new();
3191        let instrument_id = InstrumentId::from("1.234567-999004-0.0.BETFAIR");
3192        let account_id = AccountId::from("BETFAIR-001");
3193        let currency = Currency::from("GBP");
3194        let ts = UnixNanos::default();
3195
3196        // First fill: 10 @ avp=5.0
3197        let uo1 = make_test_uo(
3198            "999004",
3199            Decimal::new(20, 0),
3200            Some(Decimal::new(10, 0)),
3201            Some(Decimal::new(50, 1)),
3202        );
3203        tracker
3204            .maybe_fill_report(&uo1, uo1.s, instrument_id, account_id, currency, ts, ts)
3205            .expect("first fill");
3206
3207        // Second fill: sm=15, avp=1.0
3208        // Back-calc: (1.0*15 - 5.0*10) / 5 = (15-50)/5 = -7.0
3209        // Negative price should fall back to avp=1.0
3210        let uo2 = make_test_uo(
3211            "999004",
3212            Decimal::new(20, 0),
3213            Some(Decimal::new(15, 0)),
3214            Some(Decimal::new(10, 1)),
3215        );
3216        let fill2 = tracker
3217            .maybe_fill_report(&uo2, uo2.s, instrument_id, account_id, currency, ts, ts)
3218            .expect("second fill should use avp fallback");
3219        assert_eq!(fill2.last_qty.as_f64(), 5.0);
3220        assert_eq!(fill2.last_px.as_f64(), 1.0);
3221    }
3222
3223    #[rstest]
3224    fn test_fill_tracker_prune_clears_state() {
3225        let mut tracker = FillTracker::new();
3226        let instrument_id = InstrumentId::from("1.234567-999005-0.0.BETFAIR");
3227        let account_id = AccountId::from("BETFAIR-001");
3228        let currency = Currency::from("GBP");
3229        let ts = UnixNanos::default();
3230
3231        // Fill order fully
3232        let uo = make_test_uo(
3233            "999005",
3234            Decimal::new(10, 0),
3235            Some(Decimal::new(10, 0)),
3236            Some(Decimal::new(25, 1)),
3237        );
3238        let fill1 =
3239            tracker.maybe_fill_report(&uo, uo.s, instrument_id, account_id, currency, ts, ts);
3240        assert!(fill1.is_some());
3241
3242        // Same data again - deduplicated
3243        let fill2 =
3244            tracker.maybe_fill_report(&uo, uo.s, instrument_id, account_id, currency, ts, ts);
3245        assert!(fill2.is_none(), "should be deduplicated");
3246
3247        // Prune the bet
3248        tracker.prune("999005");
3249
3250        // After prune, same data can produce a fill again (simulates re-processing)
3251        let fill3 =
3252            tracker.maybe_fill_report(&uo, uo.s, instrument_id, account_id, currency, ts, ts);
3253        assert!(fill3.is_some(), "after prune, should produce fill again");
3254    }
3255
3256    #[rstest]
3257    fn test_fill_tracker_sm_none_returns_none() {
3258        let mut tracker = FillTracker::new();
3259
3260        // sm=None (no matched quantity field at all)
3261        let uo = make_test_uo("999006", Decimal::new(10, 0), None, None);
3262
3263        let instrument_id = InstrumentId::from("1.234567-999006-0.0.BETFAIR");
3264        let result = tracker.maybe_fill_report(
3265            &uo,
3266            uo.s,
3267            instrument_id,
3268            AccountId::from("BETFAIR-001"),
3269            Currency::from("GBP"),
3270            UnixNanos::default(),
3271            UnixNanos::default(),
3272        );
3273        assert!(result.is_none(), "None sm should not produce a fill");
3274    }
3275
3276    #[rstest]
3277    fn test_parse_order_status_report_missing_persistence_type_for_market_on_close() {
3278        let uo = UnmatchedOrder {
3279            s: Decimal::ZERO,
3280            pt: None,
3281            ot: StreamingOrderType::MarketOnClose,
3282            sr: Some(Decimal::new(10, 0)),
3283            ..make_test_uo("999007", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3284        };
3285
3286        let report = parse_order_status_report(
3287            &uo,
3288            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3289            AccountId::from("BETFAIR-001"),
3290            UnixNanos::default(),
3291            UnixNanos::default(),
3292        )
3293        .unwrap();
3294
3295        assert_eq!(report.quantity, Quantity::from("10.00"));
3296        assert_eq!(report.time_in_force, TimeInForce::AtTheClose);
3297    }
3298
3299    #[rstest]
3300    fn test_parse_order_status_report_missing_persistence_type_for_limit_on_close() {
3301        let uo = UnmatchedOrder {
3302            s: Decimal::ZERO,
3303            pt: None,
3304            ot: StreamingOrderType::LimitOnClose,
3305            sr: Some(Decimal::new(10, 0)),
3306            ..make_test_uo("999013", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3307        };
3308
3309        let report = parse_order_status_report(
3310            &uo,
3311            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3312            AccountId::from("BETFAIR-001"),
3313            UnixNanos::default(),
3314            UnixNanos::default(),
3315        )
3316        .unwrap();
3317
3318        assert_eq!(report.quantity, Quantity::from("10.00"));
3319        assert_eq!(report.time_in_force, TimeInForce::AtTheClose);
3320    }
3321
3322    #[rstest]
3323    fn test_parse_order_status_report_market_on_close_uses_bsp_liability() {
3324        let uo = UnmatchedOrder {
3325            s: Decimal::ZERO,
3326            bsp: Some(Decimal::new(20, 1)),
3327            pt: None,
3328            ot: StreamingOrderType::MarketOnClose,
3329            ..make_test_uo("999010", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3330        };
3331
3332        let report = parse_order_status_report(
3333            &uo,
3334            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3335            AccountId::from("BETFAIR-001"),
3336            UnixNanos::default(),
3337            UnixNanos::default(),
3338        )
3339        .unwrap();
3340
3341        assert_eq!(report.quantity, Quantity::from("2.00"));
3342        assert_eq!(report.time_in_force, TimeInForce::AtTheClose);
3343    }
3344
3345    #[rstest]
3346    fn test_parse_order_status_report_fails_for_non_positive_quantity() {
3347        let uo = UnmatchedOrder {
3348            s: Decimal::ZERO,
3349            sm: Some(Decimal::ZERO),
3350            sr: Some(Decimal::ZERO),
3351            sc: Some(Decimal::ZERO),
3352            sl: Some(Decimal::ZERO),
3353            sv: Some(Decimal::ZERO),
3354            ..make_test_uo("999014", Decimal::ZERO, Some(Decimal::ZERO), None)
3355        };
3356
3357        let result = parse_order_status_report(
3358            &uo,
3359            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3360            AccountId::from("BETFAIR-001"),
3361            UnixNanos::default(),
3362            UnixNanos::default(),
3363        );
3364
3365        assert!(result.is_err());
3366        assert!(
3367            result
3368                .unwrap_err()
3369                .to_string()
3370                .contains("failed to resolve positive quantity for stream order update 999014")
3371        );
3372    }
3373
3374    #[rstest]
3375    fn test_parse_order_status_report_includes_lapse_reason() {
3376        let uo = UnmatchedOrder {
3377            status: StreamingOrderStatus::ExecutionComplete,
3378            sl: Some(Decimal::ONE),
3379            lsrc: Some(crate::common::enums::LapseStatusReasonCode::SpInPlay),
3380            ..make_test_uo("999012", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3381        };
3382
3383        let report = parse_order_status_report(
3384            &uo,
3385            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3386            AccountId::from("BETFAIR-001"),
3387            UnixNanos::default(),
3388            UnixNanos::default(),
3389        )
3390        .unwrap();
3391
3392        assert_eq!(report.order_status, OrderStatus::Canceled);
3393        assert_eq!(report.cancel_reason.as_deref(), Some("SP_IN_PLAY"));
3394    }
3395
3396    #[rstest]
3397    fn test_parse_order_status_report_blank_rfo_normalizes_to_none() {
3398        // Some venues send rfo:"" instead of omitting it. ClientOrderId rejects
3399        // empty strings, so the parser must treat blank refs as missing.
3400        let uo = UnmatchedOrder {
3401            rfo: Some(String::new()),
3402            ..make_test_uo("999013", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3403        };
3404
3405        let report = parse_order_status_report(
3406            &uo,
3407            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3408            AccountId::from("BETFAIR-001"),
3409            UnixNanos::default(),
3410            UnixNanos::default(),
3411        )
3412        .unwrap();
3413
3414        assert!(report.client_order_id.is_none());
3415    }
3416
3417    #[rstest]
3418    fn test_fill_tracker_blank_rfo_normalizes_to_none() {
3419        let mut tracker = FillTracker::new();
3420        let uo = UnmatchedOrder {
3421            rfo: Some(String::new()),
3422            sm: Some(Decimal::new(5, 0)),
3423            avp: Some(Decimal::new(25, 1)),
3424            ..make_test_uo("999015", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3425        };
3426
3427        let fill = tracker
3428            .maybe_fill_report(
3429                &uo,
3430                uo.s,
3431                InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3432                AccountId::from("BETFAIR-001"),
3433                Currency::from("GBP"),
3434                UnixNanos::default(),
3435                UnixNanos::default(),
3436            )
3437            .expect("blank rfo should still produce a fill");
3438
3439        assert!(fill.client_order_id.is_none());
3440    }
3441
3442    #[rstest]
3443    fn test_fill_tracker_first_seen_void_records_without_correction() {
3444        let mut tracker = FillTracker::new();
3445        let uo = UnmatchedOrder {
3446            status: StreamingOrderStatus::ExecutionComplete,
3447            sm: Some(Decimal::new(50, 0)),
3448            sv: Some(Decimal::new(50, 0)),
3449            ..make_test_uo("999016", Decimal::new(50, 0), None, None)
3450        };
3451
3452        let first = tracker.maybe_fill_voids(&uo);
3453        let duplicate = tracker.maybe_fill_voids(&uo);
3454
3455        assert!(first.is_empty());
3456        assert!(duplicate.is_empty());
3457        assert!(!tracker.has_unseen_fill_void(&uo));
3458    }
3459
3460    #[rstest]
3461    fn test_fill_tracker_emits_only_increased_cumulative_void() {
3462        let mut tracker = FillTracker::new();
3463        let instrument_id = InstrumentId::from("1.234567-123456-0.0.BETFAIR");
3464        let account_id = AccountId::from("BETFAIR-001");
3465        let mut uo = make_test_uo(
3466            "999017",
3467            Decimal::new(60, 0),
3468            Some(Decimal::new(60, 0)),
3469            Some(Decimal::new(25, 1)),
3470        );
3471        tracker
3472            .maybe_fill_report(
3473                &uo,
3474                uo.s,
3475                instrument_id,
3476                account_id,
3477                Currency::GBP(),
3478                UnixNanos::default(),
3479                UnixNanos::default(),
3480            )
3481            .expect("matched size should establish a fill lot");
3482
3483        uo.sv = Some(Decimal::new(40, 0));
3484        let first = tracker.maybe_fill_voids(&uo);
3485        let duplicate = tracker.maybe_fill_voids(&uo);
3486        uo.sv = Some(Decimal::new(50, 0));
3487        let unseen_before_update = tracker.has_unseen_fill_void(&uo);
3488        let increased = tracker.maybe_fill_voids(&uo);
3489
3490        assert_eq!(first.len(), 1);
3491        assert_eq!(first[0].voided_qty, Quantity::from("40.00"));
3492        assert!(duplicate.is_empty());
3493        assert!(unseen_before_update);
3494        assert_eq!(increased.len(), 1);
3495        assert_eq!(increased[0].trade_id, first[0].trade_id);
3496        assert_eq!(increased[0].voided_qty, Quantity::from("50.00"));
3497        assert!(!tracker.has_unseen_fill_void(&uo));
3498    }
3499
3500    #[rstest]
3501    fn test_fill_tracker_allocates_cumulative_void_to_newest_fill_lot_first() {
3502        let mut tracker = FillTracker::new();
3503        let bet_id = "999019";
3504        let older_trade_id = TradeId::from("999019-30.00");
3505        let newer_trade_id = TradeId::from("999019-100.00");
3506        tracker.sync_fill_lot(
3507            bet_id,
3508            older_trade_id,
3509            Decimal::new(30, 0),
3510            Price::from("2.00"),
3511            Decimal::ZERO,
3512        );
3513        tracker.sync_fill_lot(
3514            bet_id,
3515            newer_trade_id,
3516            Decimal::new(70, 0),
3517            Price::from("3.00"),
3518            Decimal::ZERO,
3519        );
3520        let uo = UnmatchedOrder {
3521            sv: Some(Decimal::new(50, 0)),
3522            ..make_test_uo(bet_id, Decimal::new(100, 0), None, None)
3523        };
3524
3525        let corrections = tracker.maybe_fill_voids(&uo);
3526
3527        assert_eq!(corrections.len(), 1);
3528        assert_eq!(corrections[0].trade_id, newer_trade_id);
3529        assert_eq!(corrections[0].voided_qty, Quantity::from("50.00"));
3530        assert_eq!(corrections[0].last_px, Price::from("3.00"));
3531    }
3532
3533    #[rstest]
3534    fn test_fill_tracker_prices_new_fill_from_post_void_average() {
3535        let mut tracker = FillTracker::new();
3536        let bet_id = "999021";
3537        let older_trade_id = TradeId::from("999021-50.00");
3538        tracker.sync_fill_lot(
3539            bet_id,
3540            older_trade_id,
3541            Decimal::new(50, 0),
3542            Price::from("2.00"),
3543            Decimal::ZERO,
3544        );
3545        let (_, last_qty, last_px) = tracker
3546            .advance_cumulative_fill_with_voids(
3547                bet_id,
3548                Decimal::new(100, 0),
3549                Decimal::new(40, 0),
3550                Some(Decimal::new(25, 1)),
3551                Decimal::new(20, 1),
3552            )
3553            .expect("gross lifecycle increase should produce a fill");
3554
3555        let uo = UnmatchedOrder {
3556            sm: Some(Decimal::new(60, 0)),
3557            sv: Some(Decimal::new(40, 0)),
3558            ..make_test_uo(bet_id, Decimal::new(100, 0), None, None)
3559        };
3560        let corrections = tracker.maybe_fill_voids(&uo);
3561
3562        assert_eq!(last_qty, Quantity::from("50.00"));
3563        assert_eq!(last_px, Price::from("5.00"));
3564        assert_eq!(corrections.len(), 1);
3565        assert_eq!(corrections[0].voided_qty, Quantity::from("40.00"));
3566        assert_eq!(corrections[0].last_px, Price::from("5.00"));
3567    }
3568
3569    #[rstest]
3570    fn test_fill_tracker_preserves_prior_void_allocation_when_new_fill_arrives() {
3571        let mut tracker = FillTracker::new();
3572        let bet_id = "999020";
3573        let older_trade_id = TradeId::from("999020-60.00");
3574        let newer_trade_id = TradeId::from("999020-80.00");
3575        tracker.sync_fill_lot(
3576            bet_id,
3577            older_trade_id,
3578            Decimal::new(60, 0),
3579            Price::from("2.00"),
3580            Decimal::ZERO,
3581        );
3582        let mut uo = UnmatchedOrder {
3583            sv: Some(Decimal::new(40, 0)),
3584            ..make_test_uo(bet_id, Decimal::new(80, 0), None, None)
3585        };
3586
3587        let first = tracker.maybe_fill_voids(&uo);
3588        tracker.sync_fill_lot(
3589            bet_id,
3590            newer_trade_id,
3591            Decimal::new(20, 0),
3592            Price::from("3.00"),
3593            Decimal::ZERO,
3594        );
3595        uo.sv = Some(Decimal::new(50, 0));
3596        let increased = tracker.maybe_fill_voids(&uo);
3597
3598        assert_eq!(first.len(), 1);
3599        assert_eq!(first[0].trade_id, older_trade_id);
3600        assert_eq!(first[0].voided_qty, Quantity::from("40.00"));
3601        assert_eq!(increased.len(), 1);
3602        assert_eq!(increased[0].trade_id, newer_trade_id);
3603        assert_eq!(increased[0].voided_qty, Quantity::from("10.00"));
3604    }
3605
3606    #[rstest]
3607    fn test_fill_tracker_reconnect_increased_void_does_not_replay_gross_fill() {
3608        let mut tracker = FillTracker::new();
3609        let instrument_id = InstrumentId::from("1.234567-123456-0.0.BETFAIR");
3610        let account_id = AccountId::from("BETFAIR-001");
3611        let mut uo = make_test_uo(
3612            "999018",
3613            Decimal::new(60, 0),
3614            Some(Decimal::new(60, 0)),
3615            Some(Decimal::new(25, 1)),
3616        );
3617        let trade_id = make_trade_id(&uo);
3618        tracker.sync_fill_lot(
3619            &uo.id,
3620            trade_id,
3621            Decimal::new(60, 0),
3622            Price::from("2.50"),
3623            Decimal::new(40, 0),
3624        );
3625        tracker.sync_voided_qty(&uo.id, Decimal::new(40, 0));
3626        uo.sv = Some(Decimal::new(50, 0));
3627
3628        let fill = tracker.maybe_fill_report(
3629            &uo,
3630            uo.s,
3631            instrument_id,
3632            account_id,
3633            Currency::GBP(),
3634            UnixNanos::default(),
3635            UnixNanos::default(),
3636        );
3637        let corrections = tracker.maybe_fill_voids(&uo);
3638
3639        assert!(fill.is_none());
3640        assert_eq!(corrections.len(), 1);
3641        assert_eq!(corrections[0].trade_id, trade_id);
3642        assert_eq!(corrections[0].voided_qty, Quantity::from("50.00"));
3643    }
3644
3645    #[rstest]
3646    fn test_parse_order_status_report_missing_persistence_type_fails_for_limit_order() {
3647        let uo = UnmatchedOrder {
3648            pt: None,
3649            ..make_test_uo("999008", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3650        };
3651
3652        let result = parse_order_status_report(
3653            &uo,
3654            InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3655            AccountId::from("BETFAIR-001"),
3656            UnixNanos::default(),
3657            UnixNanos::default(),
3658        );
3659
3660        assert!(result.is_err());
3661        assert_eq!(
3662            result.unwrap_err().to_string(),
3663            "missing persistence type for order update 999008"
3664        );
3665    }
3666
3667    #[rstest]
3668    fn test_fill_tracker_uses_lifecycle_quantity_when_stream_size_is_zero() {
3669        let mut tracker = FillTracker::new();
3670        let uo = UnmatchedOrder {
3671            s: Decimal::ZERO,
3672            sr: Some(Decimal::new(10, 0)),
3673            sm: Some(Decimal::new(5, 0)),
3674            avp: Some(Decimal::new(20, 1)),
3675            ..make_test_uo("999009", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3676        };
3677
3678        let fill = tracker
3679            .maybe_fill_report(
3680                &uo,
3681                uo.s,
3682                InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3683                AccountId::from("BETFAIR-001"),
3684                Currency::from("GBP"),
3685                UnixNanos::default(),
3686                UnixNanos::default(),
3687            )
3688            .expect("zero stream size should fall back to lifecycle quantities");
3689
3690        assert_eq!(fill.last_qty, Quantity::from("5.00"));
3691    }
3692
3693    #[rstest]
3694    fn test_fill_tracker_uses_bsp_liability_when_stream_size_is_zero() {
3695        let mut tracker = FillTracker::new();
3696        let uo = UnmatchedOrder {
3697            s: Decimal::ZERO,
3698            bsp: Some(Decimal::new(20, 1)),
3699            pt: None,
3700            ot: StreamingOrderType::MarketOnClose,
3701            sm: Some(Decimal::new(10, 1)),
3702            avp: Some(Decimal::new(20, 1)),
3703            ..make_test_uo("999011", Decimal::new(10, 0), Some(Decimal::ZERO), None)
3704        };
3705
3706        let fill = tracker
3707            .maybe_fill_report(
3708                &uo,
3709                uo.s,
3710                InstrumentId::from("1.234567-123456-0.0.BETFAIR"),
3711                AccountId::from("BETFAIR-001"),
3712                Currency::from("GBP"),
3713                UnixNanos::default(),
3714                UnixNanos::default(),
3715            )
3716            .expect("zero stream size should fall back to bsp liability");
3717
3718        assert_eq!(fill.last_qty, Quantity::from("1.00"));
3719    }
3720
3721    #[rstest]
3722    fn test_fill_tracker_partial_void_still_emits_fill() {
3723        let data = load_test_json("stream/ocm_VOIDED_partial.json");
3724        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
3725
3726        if let StreamMessage::OrderChange(ocm) = msg {
3727            let oc = ocm.oc.as_ref().unwrap();
3728            let omc = &oc[0];
3729            let orc = &omc.orc.as_ref().unwrap()[0];
3730            let uo = &orc.uo.as_ref().unwrap()[0];
3731            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
3732            let ts = parse_millis_timestamp(ocm.pt);
3733
3734            let mut tracker = FillTracker::new();
3735            let fill = tracker
3736                .maybe_fill_report(
3737                    uo,
3738                    uo.s,
3739                    instrument_id,
3740                    AccountId::from("BETFAIR-001"),
3741                    Currency::GBP(),
3742                    ts,
3743                    ts,
3744                )
3745                .expect("should produce fill for matched portion");
3746
3747            // Order: s=100, sm=60, sv=40 -> fill qty should be 60
3748            assert_eq!(fill.last_qty.as_f64(), 60.0);
3749            assert_eq!(fill.last_px.as_f64(), 1.5);
3750            assert_eq!(fill.order_side, OrderSide::Sell);
3751
3752            let report = parse_order_status_report(
3753                uo,
3754                instrument_id,
3755                AccountId::from("BETFAIR-001"),
3756                ts,
3757                ts,
3758            )
3759            .unwrap();
3760            assert_eq!(report.filled_qty, Quantity::from("60.00"));
3761        } else {
3762            panic!("expected OrderChange");
3763        }
3764    }
3765
3766    #[rstest]
3767    fn test_fill_tracker_no_fill_when_sv_zero_and_fully_filled() {
3768        let data = load_test_json("stream/ocm_FILLED_sv_zero.json");
3769        let msg: StreamMessage = serde_json::from_str(&data).unwrap();
3770
3771        if let StreamMessage::OrderChange(ocm) = msg {
3772            let oc = ocm.oc.as_ref().unwrap();
3773            let omc = &oc[0];
3774            let orc = &omc.orc.as_ref().unwrap()[0];
3775            let uo = &orc.uo.as_ref().unwrap()[0];
3776            let instrument_id = make_instrument_id(&omc.id, orc.id, Decimal::ZERO);
3777            let ts = parse_millis_timestamp(ocm.pt);
3778
3779            let mut tracker = FillTracker::new();
3780            let fill = tracker
3781                .maybe_fill_report(
3782                    uo,
3783                    uo.s,
3784                    instrument_id,
3785                    AccountId::from("BETFAIR-001"),
3786                    Currency::GBP(),
3787                    ts,
3788                    ts,
3789                )
3790                .expect("fully filled order should produce fill");
3791
3792            assert_eq!(fill.last_qty.as_f64(), 50.0);
3793            assert_eq!(fill.last_px.as_f64(), 2.0);
3794
3795            // sv=0, so no void event should be generated (tested separately)
3796            assert_eq!(uo.sv, Some(Decimal::ZERO));
3797        } else {
3798            panic!("expected OrderChange");
3799        }
3800    }
3801}