Skip to main content

nautilus_backtest/
data_iterator.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Multi-stream, time-ordered data iterator for replaying historical data.
17
18use 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// TODO: block_number/transaction_index/log_index/phase are DeFi-only (zero for all other data,
29// even in non-DeFi builds); they exist to order same-block DeFi events in canonical chain order.
30// This leaks DeFi-specific shape into a general key, so it could be cfg-gated or moved behind an
31// opaque secondary key later (non-breaking, no correctness or perf cost).
32#[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
119// `as_ref` takes a closure because `DataRef`'s lifetime is an enum parameter, so its constructors
120// cannot coerce to this higher-ranked fn pointer.
121fn 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/// Internal convenience struct to keep heap entries ordered by replay key and priority.
128#[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        // min-heap on replay key, then priority sign (+/-) then index
138        self.key
139            .cmp(&other.key)
140            .then_with(|| self.priority.cmp(&other.priority))
141            .then_with(|| self.index.cmp(&other.index))
142            .reverse() // BinaryHeap is max by default -> reverse for min behavior
143    }
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/// Multi-stream, time-ordered data iterator used by the backtest engine.
153#[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, // monotonically increasing counter used to assign priorities
162}
163
164impl BacktestDataIterator {
165    /// Creates a new empty [`BacktestDataIterator`].
166    #[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    /// Adds (or replaces) a named data stream.
180    ///
181    /// When `append_data` is true the stream gets lower priority on timestamp
182    /// ties; when false (prepend) it wins ties.
183    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    /// Adds (or replaces) a named typed data stream.
194    ///
195    /// Items are ordered by replay key before insertion. A batch that arrives out of order while
196    /// sharing its backing allocation with another view is copied once, so the shared allocation
197    /// is never reordered.
198    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            // Replace existing stream - remove previous traces then re-insert below.
211            *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        // Remove old state if any
219        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    /// Removes a named data stream.
230    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            // Rebuild heap sans removed priority
237            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            // Placeholder for future generator cleanup
246        }
247    }
248
249    /// Sets the cursor of a named stream to `index` (0-based).
250    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    /// Resets all stream cursors to the beginning.
258    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    /// Returns the next backtest data element without advancing the stream cursor.
266    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    /// Advances past the current backtest data element, or does nothing if the iterator is
278    /// exhausted.
279    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        // Multi-stream path using heap
297        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        // Advance cursor and push next entry
306        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    /// Returns the next backtest data element across all streams in replay order.
323    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    /// Returns the next market [`Data`] element across all streams in chronological order.
339    #[expect(clippy::should_implement_trait)]
340    pub fn next(&mut self) -> Option<Data> {
341        self.next_item()
342    }
343
344    /// Returns whether all streams have been fully consumed.
345    #[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        // Determine if we're in single-stream mode
363        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        // Both at same timestamp, FIFO order (a before b)
581        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        // "a" is appended (lower priority), "b" is prepended (higher priority)
589        it.add_data("a", vec![quote("A.B", 100)], true);
590        it.add_data("b", vec![quote("C.D", 100)], false);
591
592        // "b" (prepend) should come first despite being added second
593        let first = it.next().unwrap();
594        let second = it.next().unwrap();
595        // Prepend stream (negative priority) wins ties over append (positive)
596        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        // Verifies that a prepend stream (negative priority) wins ties
854        // even when added after multiple append streams
855        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        // All items at the same timestamp - ordering is strictly by priority
871        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        // All should be yielded (no duplicates dropped, no items lost)
885        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}