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