1use std::collections::BinaryHeap;
19
20use ahash::AHashMap;
21use nautilus_core::UnixNanos;
22use nautilus_model::data::{BatchView, Data, DataBatch, DataRef, HasTsInit};
23
24use crate::data_batch::ReplayBatch;
25#[cfg(feature = "defi")]
26use crate::defi::replay::replay_position;
27
28#[derive(Debug, Clone, Copy, Eq, PartialEq, Ord, PartialOrd)]
33struct ReplayKey {
34 ts: UnixNanos,
35 block_number: u64,
36 transaction_index: u32,
37 log_index: u32,
38 phase: u8,
39}
40
41fn replay_key(data: DataRef<'_>) -> ReplayKey {
42 let ts = data.ts_init();
43 match data {
44 DataRef::Instrument(_)
45 | DataRef::BookDelta(_)
46 | DataRef::BookDeltas(_)
47 | DataRef::BookDepth(_)
48 | DataRef::Quote(_)
49 | DataRef::Trade(_)
50 | DataRef::Bar(_)
51 | DataRef::MarkPrice(_)
52 | DataRef::IndexPrice(_)
53 | DataRef::FundingRate(_)
54 | DataRef::OptionGreeks(_)
55 | DataRef::InstrumentStatus(_)
56 | DataRef::InstrumentClose(_)
57 | DataRef::Custom(_) => ReplayKey {
58 ts,
59 block_number: 0,
60 transaction_index: 0,
61 log_index: 0,
62 phase: 0,
63 },
64 #[cfg(feature = "defi")]
65 DataRef::Defi(defi) => {
66 let (block_number, transaction_index, log_index, phase) = replay_position(defi);
67 ReplayKey {
68 ts,
69 block_number,
70 transaction_index,
71 log_index,
72 phase,
73 }
74 }
75 }
76}
77
78fn sort_by_replay_key(batch: &mut DataBatch) {
79 match batch {
80 DataBatch::Instrument(data) => {
81 sort_view_by_replay_key(data, |item| DataRef::Instrument(item));
82 }
83 DataBatch::Custom(data) => sort_view_by_replay_key(data, |item| DataRef::Custom(item)),
84 DataBatch::BookDelta(data) => {
85 sort_view_by_replay_key(data, |item| DataRef::BookDelta(item));
86 }
87 DataBatch::BookDeltas(data) => {
88 sort_view_by_replay_key(data, |item| DataRef::BookDeltas(item));
89 }
90 DataBatch::BookDepth(data) => {
91 sort_view_by_replay_key(data, |item| DataRef::BookDepth(item));
92 }
93 DataBatch::Quote(data) => sort_view_by_replay_key(data, |item| DataRef::Quote(item)),
94 DataBatch::Trade(data) => sort_view_by_replay_key(data, |item| DataRef::Trade(item)),
95 DataBatch::Bar(data) => sort_view_by_replay_key(data, |item| DataRef::Bar(item)),
96 DataBatch::MarkPrice(data) => {
97 sort_view_by_replay_key(data, |item| DataRef::MarkPrice(item));
98 }
99 DataBatch::IndexPrice(data) => {
100 sort_view_by_replay_key(data, |item| DataRef::IndexPrice(item));
101 }
102 DataBatch::FundingRate(data) => {
103 sort_view_by_replay_key(data, |item| DataRef::FundingRate(item));
104 }
105 DataBatch::OptionGreeks(data) => {
106 sort_view_by_replay_key(data, |item| DataRef::OptionGreeks(item));
107 }
108 DataBatch::InstrumentStatus(data) => {
109 sort_view_by_replay_key(data, |item| DataRef::InstrumentStatus(item));
110 }
111 DataBatch::InstrumentClose(data) => {
112 sort_view_by_replay_key(data, |item| DataRef::InstrumentClose(item));
113 }
114 #[cfg(feature = "defi")]
115 DataBatch::Defi(data) => sort_view_by_replay_key(data, |item| DataRef::Defi(item)),
116 }
117}
118
119fn sort_view_by_replay_key<T: Clone>(view: &mut BatchView<T>, as_ref: fn(&T) -> DataRef<'_>) {
122 if !view.is_sorted_by_key(|item| replay_key(as_ref(item))) {
123 view.make_mut().sort_by_key(|item| replay_key(as_ref(item)));
124 }
125}
126
127#[derive(Debug, Eq, PartialEq)]
129struct HeapEntry {
130 key: ReplayKey,
131 priority: i32,
132 index: usize,
133}
134
135impl Ord for HeapEntry {
136 fn cmp(&self, other: &Self) -> std::cmp::Ordering {
137 self.key
139 .cmp(&other.key)
140 .then_with(|| self.priority.cmp(&other.priority))
141 .then_with(|| self.index.cmp(&other.index))
142 .reverse() }
144}
145
146impl PartialOrd for HeapEntry {
147 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
148 Some(self.cmp(other))
149 }
150}
151
152#[derive(Debug, Default)]
154pub struct BacktestDataIterator {
155 streams: AHashMap<i32, ReplayBatch>,
156 names: AHashMap<i32, String>,
157 priorities: AHashMap<String, i32>,
158 indices: AHashMap<i32, usize>,
159 heap: BinaryHeap<HeapEntry>,
160 single_priority: Option<i32>,
161 next_priority_counter: i32, }
163
164impl BacktestDataIterator {
165 #[must_use]
167 pub fn new() -> Self {
168 Self {
169 streams: AHashMap::new(),
170 names: AHashMap::new(),
171 priorities: AHashMap::new(),
172 indices: AHashMap::new(),
173 heap: BinaryHeap::new(),
174 single_priority: None,
175 next_priority_counter: 0,
176 }
177 }
178
179 pub fn add_data(&mut self, name: &str, mut data: Vec<Data>, append_data: bool) {
184 if data.is_empty() {
185 return;
186 }
187
188 data.sort_by_key(|item| replay_key(DataRef::from(item)));
189
190 self.add_stream(name, ReplayBatch::from_data(data), append_data);
191 }
192
193 pub fn add_data_batch(&mut self, name: &str, mut data: DataBatch, append_data: bool) {
199 if data.is_empty() {
200 return;
201 }
202
203 sort_by_replay_key(&mut data);
204
205 self.add_stream(name, ReplayBatch::Typed(data), append_data);
206 }
207
208 fn add_stream(&mut self, name: &str, data: ReplayBatch, append_data: bool) {
209 let priority = if let Some(p) = self.priorities.get(name) {
210 *p
212 } else {
213 self.next_priority_counter += 1;
214 let sign = if append_data { 1 } else { -1 };
215 sign * self.next_priority_counter
216 };
217
218 self.remove_data(name, true);
220
221 self.streams.insert(priority, data);
222 self.names.insert(priority, name.to_string());
223 self.priorities.insert(name.to_string(), priority);
224 self.indices.insert(priority, 0);
225
226 self.rebuild_heap();
227 }
228
229 pub fn remove_data(&mut self, name: &str, complete_remove: bool) {
231 if let Some(priority) = self.priorities.remove(name) {
232 self.streams.remove(&priority);
233 self.indices.remove(&priority);
234 self.names.remove(&priority);
235
236 self.heap.retain(|e| e.priority != priority);
238
239 if self.heap.is_empty() {
240 self.single_priority = None;
241 }
242 }
243
244 if complete_remove {
245 }
247 }
248
249 pub fn set_index(&mut self, name: &str, index: usize) {
251 if let Some(priority) = self.priorities.get(name) {
252 self.indices.insert(*priority, index);
253 self.rebuild_heap();
254 }
255 }
256
257 pub fn reset_all_cursors(&mut self) {
259 for idx in self.indices.values_mut() {
260 *idx = 0;
261 }
262 self.rebuild_heap();
263 }
264
265 pub(crate) fn peek(&self) -> Option<DataRef<'_>> {
267 if let Some(p) = self.single_priority {
268 let data = self.streams.get(&p)?;
269 let idx = *self.indices.get(&p)?;
270 return data.get(idx);
271 }
272
273 let entry = self.heap.peek()?;
274 self.streams.get(&entry.priority)?.get(entry.index)
275 }
276
277 pub(crate) fn advance(&mut self) {
280 if let Some(p) = self.single_priority {
281 let Some(data) = self.streams.get(&p) else {
282 return;
283 };
284
285 let Some(idx) = self.indices.get_mut(&p) else {
286 return;
287 };
288
289 if *idx < data.len() {
290 *idx += 1;
291 }
292
293 return;
294 }
295
296 let Some(entry) = self.heap.pop() else {
298 return;
299 };
300
301 let Some(stream) = self.streams.get(&entry.priority) else {
302 return;
303 };
304
305 let next_index = entry.index + 1;
307 self.indices.insert(entry.priority, next_index);
308
309 if next_index < stream.len() {
310 self.heap.push(HeapEntry {
311 key: replay_key(
312 stream
313 .get(next_index)
314 .expect("next index is within the data batch"),
315 ),
316 priority: entry.priority,
317 index: next_index,
318 });
319 }
320 }
321
322 pub(crate) fn next_item(&mut self) -> Option<Data> {
324 let element = if let Some(p) = self.single_priority {
325 let data = self.streams.get(&p)?;
326 let idx = *self.indices.get(&p)?;
327 data.get_owned(idx)?
328 } else {
329 let entry = self.heap.peek()?;
330 self.streams.get(&entry.priority)?.get_owned(entry.index)?
331 };
332
333 self.advance();
334
335 Some(element)
336 }
337
338 #[expect(clippy::should_implement_trait)]
340 pub fn next(&mut self) -> Option<Data> {
341 self.next_item()
342 }
343
344 #[must_use]
346 pub fn is_done(&self) -> bool {
347 if let Some(p) = self.single_priority {
348 if let Some(idx) = self.indices.get(&p)
349 && let Some(data) = self.streams.get(&p)
350 {
351 return *idx >= data.len();
352 }
353 true
354 } else {
355 self.heap.is_empty()
356 }
357 }
358
359 fn rebuild_heap(&mut self) {
360 self.heap.clear();
361
362 if self.streams.len() == 1 {
364 self.single_priority = self.streams.keys().next().copied();
365 return;
366 }
367 self.single_priority = None;
368
369 for (&priority, data) in &self.streams {
370 let idx = *self.indices.get(&priority).unwrap_or(&0);
371 if idx < data.len() {
372 self.heap.push(HeapEntry {
373 key: replay_key(data.get(idx).expect("index is within the data batch")),
374 priority,
375 index: idx,
376 });
377 }
378 }
379 }
380}
381
382#[cfg(test)]
383mod tests {
384 use std::sync::Arc;
385
386 use nautilus_model::{
387 data::{
388 Bar, FundingRateUpdate, IndexPriceUpdate, InstrumentClose, InstrumentStatus,
389 MarkPriceUpdate, OptionGreeks, OrderBookDelta, OrderBookDeltas, OrderBookDepth,
390 QuoteTick, TradeTick,
391 stubs::{
392 stub_bar, stub_delta, stub_deltas, stub_depth10, stub_instrument_close,
393 stub_instrument_status, stub_trade_ethusdt_buy,
394 },
395 },
396 identifiers::InstrumentId,
397 types::{Price, Quantity},
398 };
399 #[cfg(feature = "defi")]
400 use nautilus_model::{
401 defi::{
402 DefiData,
403 data::block::BlockPosition,
404 pool_analysis::snapshot::{PoolAnalytics, PoolSnapshot, PoolState},
405 },
406 identifiers::{Symbol, Venue},
407 };
408 use rstest::rstest;
409
410 use super::*;
411
412 fn quote_tick(id: &str, ts: u64) -> QuoteTick {
413 QuoteTick::new(
414 InstrumentId::from(id),
415 Price::from("1.0"),
416 Price::from("1.0"),
417 Quantity::from(100),
418 Quantity::from(100),
419 ts.into(),
420 ts.into(),
421 )
422 }
423
424 fn quote(id: &str, ts: u64) -> Data {
425 Data::Quote(quote_tick(id, ts))
426 }
427
428 fn trade(id: &str, ts: u64) -> Data {
429 let mut trade = stub_trade_ethusdt_buy();
430 trade.instrument_id = InstrumentId::from(id);
431 trade.ts_event = UnixNanos::from(ts);
432 trade.ts_init = UnixNanos::from(ts);
433 Data::Trade(trade)
434 }
435
436 fn collect_sequence(it: &mut BacktestDataIterator) -> Vec<(InstrumentId, UnixNanos)> {
437 let mut sequence = Vec::new();
438 while let Some(data) = it.peek() {
439 sequence.push((data.instrument_id(), data.ts_init()));
440 it.advance();
441 }
442 sequence
443 }
444
445 fn collect_ts(it: &mut BacktestDataIterator) -> Vec<u64> {
446 let mut ts = Vec::new();
447 while let Some(d) = it.next() {
448 ts.push(d.ts_init().as_u64());
449 }
450 ts
451 }
452
453 #[cfg(feature = "defi")]
454 fn defi_pool_snapshot(ts: u64, block: u64, transaction_index: u32, log_index: u32) -> DefiData {
455 let instrument_id = InstrumentId::new(Symbol::from("ETH/USDC"), Venue::from("UNISWAPV3"));
456 let snapshot = PoolSnapshot::new(
457 instrument_id,
458 PoolState::default(),
459 Vec::new(),
460 Vec::new(),
461 PoolAnalytics::default(),
462 BlockPosition::new(block, format!("0x{block:x}"), transaction_index, log_index),
463 UnixNanos::from(ts),
464 UnixNanos::from(ts),
465 );
466
467 DefiData::PoolSnapshot(snapshot)
468 }
469
470 #[cfg(feature = "defi")]
471 fn defi_snapshot(ts: u64, block: u64, transaction_index: u32, log_index: u32) -> Data {
472 Data::Defi(Box::new(defi_pool_snapshot(
473 ts,
474 block,
475 transaction_index,
476 log_index,
477 )))
478 }
479
480 #[rstest]
481 fn test_single_stream_yields_in_order() {
482 let mut it = BacktestDataIterator::new();
483 it.add_data(
484 "s",
485 vec![quote("A.B", 100), quote("A.B", 200), quote("A.B", 300)],
486 true,
487 );
488
489 assert_eq!(collect_ts(&mut it), vec![100, 200, 300]);
490 assert!(it.is_done());
491 }
492
493 #[rstest]
494 fn test_single_stream_exhaustion_returns_none() {
495 let mut it = BacktestDataIterator::new();
496 it.add_data("s", vec![quote("A.B", 1), quote("A.B", 3)], true);
497 assert_eq!(it.next().unwrap().ts_init(), UnixNanos::from(1));
498 assert_eq!(it.next().unwrap().ts_init(), UnixNanos::from(3));
499 assert!(it.next().is_none());
500 }
501
502 #[rstest]
503 fn test_peek_does_not_consume_single_stream_item() {
504 let mut it = BacktestDataIterator::new();
505 it.add_data("s", vec![quote("A.B", 1), quote("A.B", 2)], true);
506
507 assert_eq!(it.peek().unwrap().ts_init(), UnixNanos::from(1));
508 assert_eq!(it.peek().unwrap().ts_init(), UnixNanos::from(1));
509 assert_eq!(it.next().unwrap().ts_init(), UnixNanos::from(1));
510 assert_eq!(it.peek().unwrap().ts_init(), UnixNanos::from(2));
511 }
512
513 #[rstest]
514 fn test_single_stream_sorts_unsorted_input() {
515 let mut it = BacktestDataIterator::new();
516 it.add_data(
517 "s",
518 vec![quote("A.B", 300), quote("A.B", 100), quote("A.B", 200)],
519 true,
520 );
521
522 assert_eq!(collect_ts(&mut it), vec![100, 200, 300]);
523 }
524
525 #[rstest]
526 fn test_two_stream_merge_chronological() {
527 let mut it = BacktestDataIterator::new();
528 it.add_data("s1", vec![quote("A.B", 1), quote("A.B", 4)], true);
529 it.add_data("s2", vec![quote("C.D", 2), quote("C.D", 3)], false);
530
531 assert_eq!(collect_ts(&mut it), vec![1, 2, 3, 4]);
532 }
533
534 #[rstest]
535 fn test_peek_does_not_consume_multi_stream_heap_item() {
536 let mut it = BacktestDataIterator::new();
537 it.add_data("s1", vec![quote("A.B", 1), quote("A.B", 4)], true);
538 it.add_data("s2", vec![quote("C.D", 2), quote("C.D", 3)], true);
539
540 assert_eq!(it.peek().unwrap().ts_init(), UnixNanos::from(1));
541 assert_eq!(it.peek().unwrap().ts_init(), UnixNanos::from(1));
542 assert_eq!(it.next().unwrap().ts_init(), UnixNanos::from(1));
543 assert_eq!(it.peek().unwrap().ts_init(), UnixNanos::from(2));
544 assert_eq!(collect_ts(&mut it), vec![2, 3, 4]);
545 }
546
547 #[rstest]
548 fn test_three_stream_merge_sorted() {
549 let mut it = BacktestDataIterator::new();
550 let data_len = 5;
551 let d0: Vec<Data> = (0..data_len).map(|k| quote("A.B", 3 * k)).collect();
552 let d1: Vec<Data> = (0..data_len).map(|k| quote("C.D", 3 * k + 1)).collect();
553 let d2: Vec<Data> = (0..data_len).map(|k| quote("E.F", 3 * k + 2)).collect();
554 it.add_data("d0", d0, true);
555 it.add_data("d1", d1, true);
556 it.add_data("d2", d2, true);
557
558 let ts = collect_ts(&mut it);
559 assert_eq!(ts.len(), 15);
560 for i in 0..ts.len() - 1 {
561 assert!(ts[i] <= ts[i + 1], "Not sorted at index {i}");
562 }
563 }
564
565 #[rstest]
566 fn test_multiple_streams_merge_order() {
567 let mut it = BacktestDataIterator::new();
568 it.add_data("s1", vec![quote("A.B", 100), quote("A.B", 300)], true);
569 it.add_data("s2", vec![quote("C.D", 200), quote("C.D", 400)], true);
570
571 assert_eq!(collect_ts(&mut it), vec![100, 200, 300, 400]);
572 }
573
574 #[rstest]
575 fn test_append_data_priority_default_fifo() {
576 let mut it = BacktestDataIterator::new();
577 it.add_data("a", vec![quote("A.B", 100)], true);
578 it.add_data("b", vec![quote("C.D", 100)], true);
579
580 let ts = collect_ts(&mut it);
582 assert_eq!(ts, vec![100, 100]);
583 }
584
585 #[rstest]
586 fn test_prepend_priority_wins_ties() {
587 let mut it = BacktestDataIterator::new();
588 it.add_data("a", vec![quote("A.B", 100)], true);
590 it.add_data("b", vec![quote("C.D", 100)], false);
591
592 let first = it.next().unwrap();
594 let second = it.next().unwrap();
595 assert_eq!(first.instrument_id(), InstrumentId::from("C.D"));
597 assert_eq!(second.instrument_id(), InstrumentId::from("A.B"));
598 }
599
600 #[rstest]
601 fn test_is_done_empty_iterator() {
602 let it = BacktestDataIterator::new();
603 assert!(it.is_done());
604 }
605
606 #[rstest]
607 fn test_is_done_after_consumption() {
608 let mut it = BacktestDataIterator::new();
609 it.add_data("s", vec![quote("A.B", 1)], true);
610
611 assert!(!it.is_done());
612 it.next();
613 assert!(it.is_done());
614 }
615
616 #[rstest]
617 fn test_is_done_multi_stream() {
618 let mut it = BacktestDataIterator::new();
619 it.add_data("s1", vec![quote("A.B", 1)], true);
620 it.add_data("s2", vec![quote("C.D", 2)], true);
621
622 assert!(!it.is_done());
623 it.next();
624 assert!(!it.is_done());
625 it.next();
626 assert!(it.is_done());
627 }
628
629 #[rstest]
630 fn test_partial_consumption_then_complete() {
631 let mut it = BacktestDataIterator::new();
632 it.add_data(
633 "s",
634 vec![
635 quote("A.B", 0),
636 quote("A.B", 1),
637 quote("A.B", 2),
638 quote("A.B", 3),
639 ],
640 true,
641 );
642
643 assert_eq!(it.next().unwrap().ts_init().as_u64(), 0);
644 assert_eq!(it.next().unwrap().ts_init().as_u64(), 1);
645
646 let remaining = collect_ts(&mut it);
647 assert_eq!(remaining, vec![2, 3]);
648 assert!(it.is_done());
649 }
650
651 #[rstest]
652 fn test_remove_stream_reduces_output() {
653 let mut it = BacktestDataIterator::new();
654 it.add_data("a", vec![quote("A.B", 1)], true);
655 it.add_data("b", vec![quote("C.D", 2)], true);
656
657 it.remove_data("a", false);
658
659 assert_eq!(collect_ts(&mut it), vec![2]);
660 }
661
662 #[rstest]
663 fn test_remove_all_streams_yields_empty() {
664 let mut it = BacktestDataIterator::new();
665 it.add_data("x", vec![quote("A.B", 1)], true);
666 it.add_data("y", vec![quote("C.D", 2)], true);
667
668 it.remove_data("x", false);
669 it.remove_data("y", false);
670
671 assert!(it.next().is_none());
672 assert!(it.is_done());
673 }
674
675 #[rstest]
676 fn test_remove_nonexistent_stream_is_noop() {
677 let mut it = BacktestDataIterator::new();
678 it.add_data("s", vec![quote("A.B", 1)], true);
679
680 it.remove_data("nonexistent", false);
681
682 assert_eq!(collect_ts(&mut it), vec![1]);
683 }
684
685 #[rstest]
686 fn test_remove_after_full_consumption() {
687 let mut it = BacktestDataIterator::new();
688 it.add_data("s", vec![quote("A.B", 1), quote("A.B", 2)], true);
689
690 collect_ts(&mut it);
691
692 it.remove_data("s", true);
693 assert!(it.is_done());
694 }
695
696 #[rstest]
697 fn test_set_index_rewinds_stream() {
698 let mut it = BacktestDataIterator::new();
699 it.add_data(
700 "s",
701 vec![quote("A.B", 10), quote("A.B", 20), quote("A.B", 30)],
702 true,
703 );
704
705 assert_eq!(it.next().unwrap().ts_init().as_u64(), 10);
706
707 it.set_index("s", 0);
708
709 assert_eq!(collect_ts(&mut it), vec![10, 20, 30]);
710 }
711
712 #[rstest]
713 fn test_set_index_skips_forward() {
714 let mut it = BacktestDataIterator::new();
715 it.add_data(
716 "s",
717 vec![quote("A.B", 10), quote("A.B", 20), quote("A.B", 30)],
718 true,
719 );
720
721 it.set_index("s", 2);
722
723 assert_eq!(collect_ts(&mut it), vec![30]);
724 }
725
726 #[rstest]
727 fn test_set_index_uses_logical_mixed_stream_offset() {
728 let mut it = BacktestDataIterator::new();
729 it.add_data(
730 "mixed",
731 vec![quote("A.B", 10), trade("C.D", 10), quote("E.F", 20)],
732 true,
733 );
734
735 it.set_index("mixed", 1);
736
737 assert_eq!(
738 collect_sequence(&mut it),
739 vec![
740 (InstrumentId::from("C.D"), UnixNanos::from(10)),
741 (InstrumentId::from("E.F"), UnixNanos::from(20)),
742 ]
743 );
744 }
745
746 #[rstest]
747 fn test_set_index_nonexistent_stream_is_noop() {
748 let mut it = BacktestDataIterator::new();
749 it.add_data("s", vec![quote("A.B", 1)], true);
750
751 it.set_index("nonexistent", 0);
752
753 assert_eq!(collect_ts(&mut it), vec![1]);
754 }
755
756 #[rstest]
757 fn test_reset_all_cursors_single_stream() {
758 let mut it = BacktestDataIterator::new();
759 it.add_data("s", vec![quote("A.B", 1), quote("A.B", 2)], true);
760
761 collect_ts(&mut it);
762 assert!(it.is_done());
763
764 it.reset_all_cursors();
765 assert!(!it.is_done());
766 assert_eq!(collect_ts(&mut it), vec![1, 2]);
767 }
768
769 #[rstest]
770 fn test_reset_all_cursors_multi_stream() {
771 let mut it = BacktestDataIterator::new();
772 it.add_data("s1", vec![quote("A.B", 1), quote("A.B", 3)], true);
773 it.add_data("s2", vec![quote("C.D", 2), quote("C.D", 4)], true);
774
775 collect_ts(&mut it);
776 assert!(it.is_done());
777
778 it.reset_all_cursors();
779 assert_eq!(collect_ts(&mut it), vec![1, 2, 3, 4]);
780 }
781
782 #[rstest]
783 fn test_readding_data_replaces_stream() {
784 let mut it = BacktestDataIterator::new();
785 it.add_data("X", vec![quote("A.B", 1), quote("A.B", 2)], true);
786 it.add_data("X", vec![quote("A.B", 10)], true);
787
788 assert_eq!(collect_ts(&mut it), vec![10]);
789 }
790
791 #[rstest]
792 fn test_readding_data_reuses_equal_key_stream_priority() {
793 let mut it = BacktestDataIterator::new();
794 it.add_data("first", vec![quote("A.B", 10)], true);
795 it.add_data("second", vec![quote("C.D", 10)], true);
796 it.add_data("first", vec![quote("E.F", 10)], true);
797
798 assert_eq!(
799 collect_sequence(&mut it),
800 vec![
801 (InstrumentId::from("E.F"), UnixNanos::from(10)),
802 (InstrumentId::from("C.D"), UnixNanos::from(10)),
803 ]
804 );
805 }
806
807 #[rstest]
808 fn test_add_empty_data_is_noop() {
809 let mut it = BacktestDataIterator::new();
810 it.add_data("empty", vec![], true);
811
812 assert!(it.is_done());
813 assert!(it.next().is_none());
814 }
815
816 #[rstest]
817 fn test_empty_iterator_returns_none() {
818 let mut it = BacktestDataIterator::new();
819 assert!(it.next().is_none());
820 assert!(it.is_done());
821 }
822
823 #[rstest]
824 fn test_multiple_add_data_calls_with_different_names() {
825 let mut it = BacktestDataIterator::new();
826 it.add_data("batch_0", vec![quote("A.B", 1), quote("A.B", 3)], true);
827 it.add_data("batch_1", vec![quote("A.B", 2), quote("A.B", 4)], true);
828
829 assert_eq!(collect_ts(&mut it), vec![1, 2, 3, 4]);
830 }
831
832 #[rstest]
833 fn test_typed_and_compatibility_streams_preserve_equal_key_order() {
834 let mut single_batch = BacktestDataIterator::new();
835 single_batch.add_data(
836 "single",
837 vec![quote("A.B", 10), trade("C.D", 10), quote("E.F", 20)],
838 true,
839 );
840
841 let mut split_batches = BacktestDataIterator::new();
842 split_batches.add_data("first", vec![quote("A.B", 10)], true);
843 split_batches.add_data("second", vec![trade("C.D", 10), quote("E.F", 20)], true);
844
845 assert_eq!(
846 collect_sequence(&mut split_batches),
847 collect_sequence(&mut single_batch)
848 );
849 }
850
851 #[rstest]
852 fn test_prepend_stream_always_wins_ties_across_batches() {
853 let mut it = BacktestDataIterator::new();
856 it.add_data("append_a", vec![quote("A.B", 100)], true);
857 it.add_data("append_b", vec![quote("C.D", 100)], true);
858 it.add_data("prepend", vec![quote("E.F", 100)], false);
859
860 let first = it.next().unwrap();
861 assert_eq!(
862 first.instrument_id(),
863 InstrumentId::from("E.F"),
864 "Prepend stream should always come first in ties"
865 );
866 }
867
868 #[rstest]
869 fn test_equal_timestamps_across_many_streams_preserves_priority_order() {
870 let mut it = BacktestDataIterator::new();
872 it.add_data("s1", vec![quote("A.B", 50)], true);
873 it.add_data("s2", vec![quote("C.D", 50)], true);
874 it.add_data("s3", vec![quote("E.F", 50)], true);
875 it.add_data("s4", vec![quote("G.H", 50)], true);
876
877 let mut ids = Vec::new();
878 while let Some(d) = it.next() {
879 ids.push(d.instrument_id());
880 }
881
882 assert_eq!(ids.len(), 4);
883
884 assert!(ids.contains(&InstrumentId::from("A.B")));
886 assert!(ids.contains(&InstrumentId::from("C.D")));
887 assert!(ids.contains(&InstrumentId::from("E.F")));
888 assert!(ids.contains(&InstrumentId::from("G.H")));
889 }
890
891 #[rstest]
892 fn test_add_data_batch_sorts_every_typed_family_by_replay_key() {
893 let instrument_id = InstrumentId::from("A.B");
894 let late = UnixNanos::from(200);
895 let early = UnixNanos::from(100);
896 let mark = |ts| MarkPriceUpdate::new(instrument_id, Price::from("1.0"), ts, ts);
897 let index = |ts| IndexPriceUpdate::new(instrument_id, Price::from("1.0"), ts, ts);
898 let funding = |ts| {
899 FundingRateUpdate::new(instrument_id, "0.0001".parse().unwrap(), None, None, ts, ts)
900 };
901 let greeks = |ts| OptionGreeks {
902 instrument_id,
903 ts_init: ts,
904 ..OptionGreeks::default()
905 };
906 let batches = vec![
907 DataBatch::from(vec![
908 OrderBookDelta {
909 ts_init: late,
910 ..stub_delta()
911 },
912 OrderBookDelta {
913 ts_init: early,
914 ..stub_delta()
915 },
916 ]),
917 DataBatch::from(vec![
918 OrderBookDeltas {
919 ts_init: late,
920 ..stub_deltas()
921 },
922 OrderBookDeltas {
923 ts_init: early,
924 ..stub_deltas()
925 },
926 ]),
927 DataBatch::from(vec![
928 OrderBookDepth {
929 ts_init: late,
930 ..stub_depth10()
931 },
932 OrderBookDepth {
933 ts_init: early,
934 ..stub_depth10()
935 },
936 ]),
937 DataBatch::from(vec![quote_tick("A.B", 200), quote_tick("A.B", 100)]),
938 DataBatch::from(vec![
939 TradeTick {
940 ts_init: late,
941 ..stub_trade_ethusdt_buy()
942 },
943 TradeTick {
944 ts_init: early,
945 ..stub_trade_ethusdt_buy()
946 },
947 ]),
948 DataBatch::from(vec![
949 Bar {
950 ts_init: late,
951 ..stub_bar()
952 },
953 Bar {
954 ts_init: early,
955 ..stub_bar()
956 },
957 ]),
958 DataBatch::from(vec![mark(late), mark(early)]),
959 DataBatch::from(vec![index(late), index(early)]),
960 DataBatch::from(vec![funding(late), funding(early)]),
961 DataBatch::from(vec![greeks(late), greeks(early)]),
962 DataBatch::from(vec![
963 InstrumentStatus {
964 ts_init: late,
965 ..stub_instrument_status()
966 },
967 InstrumentStatus {
968 ts_init: early,
969 ..stub_instrument_status()
970 },
971 ]),
972 DataBatch::from(vec![
973 InstrumentClose {
974 ts_init: late,
975 ..stub_instrument_close()
976 },
977 InstrumentClose {
978 ts_init: early,
979 ..stub_instrument_close()
980 },
981 ]),
982 ];
983 assert_eq!(
984 batches.len(),
985 12,
986 "every static DataBatch variant needs a case"
987 );
988
989 for batch in batches {
990 let mut it = BacktestDataIterator::new();
991 it.add_data_batch("typed", batch, true);
992
993 assert_eq!(collect_ts(&mut it), vec![100, 200]);
994 }
995 }
996
997 #[rstest]
998 fn test_add_data_batch_shares_sorted_backing_allocation() {
999 let data = Arc::new(vec![quote_tick("A.B", 100), quote_tick("A.B", 200)]);
1000 let mut it = BacktestDataIterator::new();
1001
1002 it.add_data_batch(
1003 "typed",
1004 DataBatch::Quote(BatchView::from(Arc::clone(&data))),
1005 true,
1006 );
1007
1008 assert_eq!(Arc::strong_count(&data), 2);
1009 assert_eq!(collect_ts(&mut it), vec![100, 200]);
1010 }
1011
1012 #[rstest]
1013 fn test_add_data_batch_copies_shared_unsorted_backing_allocation() {
1014 let data = Arc::new(vec![quote_tick("A.B", 200), quote_tick("A.B", 100)]);
1015 let mut it = BacktestDataIterator::new();
1016
1017 it.add_data_batch(
1018 "typed",
1019 DataBatch::Quote(BatchView::from(Arc::clone(&data))),
1020 true,
1021 );
1022
1023 assert_eq!(Arc::strong_count(&data), 1);
1024 assert_eq!(data[0].ts_init, UnixNanos::from(200));
1025 assert_eq!(collect_ts(&mut it), vec![100, 200]);
1026 }
1027
1028 #[rstest]
1029 fn test_add_empty_data_batch_is_noop() {
1030 let mut it = BacktestDataIterator::new();
1031 it.add_data_batch("typed", DataBatch::from(Vec::<QuoteTick>::new()), true);
1032
1033 assert!(it.is_done());
1034 assert!(it.peek().is_none());
1035 }
1036
1037 #[rstest]
1038 fn test_typed_batch_and_legacy_streams_preserve_equal_key_order() {
1039 let mut single_batch = BacktestDataIterator::new();
1040 single_batch.add_data(
1041 "single",
1042 vec![quote("A.B", 10), trade("C.D", 10), quote("E.F", 20)],
1043 true,
1044 );
1045
1046 let mut split_streams = BacktestDataIterator::new();
1047 split_streams.add_data_batch(
1048 "typed",
1049 DataBatch::from(vec![quote_tick("A.B", 10), quote_tick("E.F", 20)]),
1050 true,
1051 );
1052 split_streams.add_data("legacy", vec![trade("C.D", 10)], true);
1053
1054 assert_eq!(
1055 collect_sequence(&mut split_streams),
1056 collect_sequence(&mut single_batch)
1057 );
1058 }
1059
1060 #[cfg(feature = "defi")]
1061 #[rstest]
1062 fn test_add_data_batch_orders_defi_by_block_position() {
1063 let mut it = BacktestDataIterator::new();
1064 it.add_data_batch(
1065 "defi",
1066 DataBatch::from(vec![
1067 defi_pool_snapshot(100, 12, 4, 1),
1068 defi_pool_snapshot(100, 11, 9, 9),
1069 defi_pool_snapshot(100, 12, 2, 7),
1070 ]),
1071 true,
1072 );
1073
1074 let mut positions = Vec::new();
1075 while let Some(Data::Defi(data)) = it.next_item() {
1076 positions.push(data.block_position());
1077 }
1078
1079 assert_eq!(positions, vec![(11, 9, 9), (12, 2, 7), (12, 4, 1)]);
1080 }
1081
1082 #[cfg(feature = "defi")]
1083 #[rstest]
1084 fn test_defi_data_orders_equal_timestamps_by_block_position() {
1085 let mut it = BacktestDataIterator::new();
1086 it.add_data(
1087 "defi",
1088 vec![
1089 defi_snapshot(100, 12, 4, 1),
1090 defi_snapshot(100, 11, 9, 9),
1091 defi_snapshot(100, 12, 2, 7),
1092 ],
1093 true,
1094 );
1095
1096 let mut positions = Vec::new();
1097 while let Some(Data::Defi(data)) = it.next_item() {
1098 positions.push(data.block_position());
1099 }
1100
1101 assert_eq!(positions, vec![(11, 9, 9), (12, 2, 7), (12, 4, 1)]);
1102 }
1103}