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