1use 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
56pub 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 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#[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#[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 (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
242pub 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#[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#[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 #[must_use]
285 pub fn new() -> Self {
286 Self::default()
287 }
288
289 #[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 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 #[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 #[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 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 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 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 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#[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#[must_use]
775pub fn is_lapsed(uo: &UnmatchedOrder) -> bool {
776 uo.status == StreamingOrderStatus::ExecutionComplete && uo.lsrc.is_some()
777}
778
779pub 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 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 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#[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#[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#[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#[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#[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#[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#[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 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 assert_eq!(deltas.deltas[0].action, BookAction::Clear);
1245 assert!(RecordFlag::F_SNAPSHOT.matches(deltas.deltas[0].flags));
1246
1247 for delta in &deltas.deltas[1..] {
1249 assert_eq!(delta.action, BookAction::Add);
1250 assert!(RecordFlag::F_SNAPSHOT.matches(delta.flags));
1251 }
1252
1253 let last = deltas.deltas.last().unwrap();
1255 assert!(RecordFlag::F_LAST.matches(last.flags));
1256
1257 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 assert!(deltas.deltas.iter().all(|d| d.action != BookAction::Clear));
1333
1334 let last = deltas.deltas.last().unwrap();
1336 assert!(RecordFlag::F_LAST.matches(last.flags));
1337
1338 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 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 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 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 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 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)), 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 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 assert_eq!(report.order_status, OrderStatus::PartiallyFilled);
1645 assert_eq!(report.order_side, Some(OrderSide::Sell)); 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)); assert_eq!(report.filled_qty.as_f64(), 10.0);
1680 assert_eq!(report.quantity.as_f64(), 10.0);
1681
1682 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)); 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 assert_eq!(report.order_status, OrderStatus::PartiallyFilled);
1742 assert_eq!(report.order_side, Some(OrderSide::Buy)); 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 assert_eq!(report.order_status, OrderStatus::PartiallyFilled);
1772 assert_eq!(report.order_side, Some(OrderSide::Sell)); 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 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 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 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 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 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 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 let fill1 =
2028 tracker.maybe_fill_report(uo, uo.s, instrument_id, account_id, currency, ts, ts);
2029 assert!(fill1.is_some());
2030
2031 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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.sync_order("123456", Decimal::new(10, 0), Decimal::new(25, 1));
2773
2774 tracker.sync_order("123456", Decimal::new(5, 0), Decimal::new(20, 1));
2776
2777 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 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 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 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 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 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 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 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 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 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 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 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 tracker.prune("999005");
3120
3121 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 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 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 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 assert_eq!(uo.sv, Some(Decimal::ZERO));
3668 } else {
3669 panic!("expected OrderChange");
3670 }
3671 }
3672}