Skip to main content

qs_backtest/
data_feed.rs

1//! Data feed abstraction for backtesting.
2//!
3//! A [`DataFeed`] produces a time-ordered sequence of [`MarketEvent`]s that the
4//! backtest runner consumes.  Two built-in implementations are provided:
5//!
6//! - [`VecFeed`] — wraps a pre-loaded `Vec<MarketEvent>` (useful for tests and
7//!   any data source you can materialise up-front).
8//! - Conversion helpers from `qs-data-preprocess` types (`Tick`, `Bar`) so
9//!   you can query DuckDB, convert the results, and feed them straight in.
10
11use std::cmp::Reverse;
12use std::collections::BinaryHeap;
13use std::marker::PhantomData;
14
15use chrono::NaiveDateTime;
16use qs_core::ExecutionPricer;
17use qs_core::types::PriceQuote;
18use serde::{Deserialize, Serialize};
19use thiserror::Error;
20
21// ─── MarketEvent ────────────────────────────────────────────────────────────
22
23/// A single market data event consumed by the backtest runner.
24#[derive(Debug, Clone, Serialize, Deserialize)]
25pub enum MarketEvent {
26    /// A bid/ask tick.
27    Tick {
28        symbol: String,
29        ts: NaiveDateTime,
30        bid: f64,
31        ask: f64,
32    },
33    /// An OHLCV bar.
34    Bar {
35        symbol: String,
36        ts: NaiveDateTime,
37        open: f64,
38        high: f64,
39        low: f64,
40        close: f64,
41        volume: i64,
42    },
43}
44
45impl MarketEvent {
46    /// Timestamp of the event.
47    pub fn ts(&self) -> NaiveDateTime {
48        match self {
49            MarketEvent::Tick { ts, .. } => *ts,
50            MarketEvent::Bar { ts, .. } => *ts,
51        }
52    }
53
54    /// Symbol of the event.
55    pub fn symbol(&self) -> &str {
56        match self {
57            MarketEvent::Tick { symbol, .. } => symbol,
58            MarketEvent::Bar { symbol, .. } => symbol,
59        }
60    }
61
62    /// Convert the event into a [`PriceQuote`] suitable for the trade engine.
63    ///
64    /// - For ticks this is straightforward (bid/ask).
65    /// - For bars the close price is used for both bid and ask (zero spread
66    ///   approximation).  A more sophisticated feed could model the spread
67    ///   separately.
68    pub fn to_quote(&self) -> PriceQuote {
69        match self {
70            MarketEvent::Tick {
71                symbol,
72                ts,
73                bid,
74                ask,
75            } => PriceQuote {
76                symbol: symbol.clone(),
77                ts: *ts,
78                bid: *bid,
79                ask: *ask,
80            },
81            MarketEvent::Bar {
82                symbol, ts, close, ..
83            } => PriceQuote {
84                symbol: symbol.clone(),
85                ts: *ts,
86                bid: *close,
87                ask: *close,
88            },
89        }
90    }
91
92    /// Convert this event into a quote only when its executable bid/ask view is
93    /// finite, positive, and not crossed.
94    pub fn to_valid_quote(&self) -> Option<PriceQuote> {
95        let quote = self.to_quote();
96        ExecutionPricer::validate_quote(&quote).ok().map(|()| quote)
97    }
98}
99
100/// Roles a feed event serves in a multi-series backtest.
101#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
102pub struct SeriesRoles {
103    pub primary: bool,
104    pub conversion: bool,
105}
106
107impl SeriesRoles {
108    pub const PRIMARY: Self = Self {
109        primary: true,
110        conversion: false,
111    };
112    pub const CONVERSION: Self = Self {
113        primary: false,
114        conversion: true,
115    };
116    pub const PRIMARY_AND_CONVERSION: Self = Self {
117        primary: true,
118        conversion: true,
119    };
120}
121
122/// Deterministic identity and ordering metadata for one market event.
123#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
124pub struct EventMetadata {
125    pub roles: SeriesRoles,
126    pub series_rank: u32,
127    /// Physical source-row ordinal when supplied by a streaming source, otherwise the emitted event ordinal.
128    pub row_sequence: u64,
129}
130
131impl EventMetadata {
132    pub const fn new(roles: SeriesRoles, series_rank: u32, row_sequence: u64) -> Self {
133        Self {
134            roles,
135            series_rank,
136            row_sequence,
137        }
138    }
139}
140
141/// A market event paired with non-persisted feed metadata.
142#[derive(Debug, Clone, Serialize, Deserialize)]
143pub struct FeedEvent {
144    pub event: MarketEvent,
145    pub metadata: EventMetadata,
146}
147
148impl FeedEvent {
149    pub fn new(event: MarketEvent, metadata: EventMetadata) -> Self {
150        Self { event, metadata }
151    }
152
153    /// Ordering key used by deterministic feeds.
154    pub fn ordering_key(&self) -> (NaiveDateTime, u32, u64) {
155        (
156            self.event.ts(),
157            self.metadata.series_rank,
158            self.metadata.row_sequence,
159        )
160    }
161}
162
163/// A streaming event paired with its physical source-row ordinal.
164#[derive(Debug, Clone)]
165pub struct SequencedMarketEvent {
166    pub event: MarketEvent,
167    pub source_row_ordinal: u64,
168}
169
170impl SequencedMarketEvent {
171    pub fn new(event: MarketEvent, source_row_ordinal: u64) -> Self {
172        Self {
173            event,
174            source_row_ordinal,
175        }
176    }
177}
178
179/// Converts an event source item into an event and optional physical ordinal.
180pub trait EventSourceItem {
181    fn into_event_and_ordinal(self) -> (MarketEvent, Option<u64>);
182}
183
184impl EventSourceItem for MarketEvent {
185    fn into_event_and_ordinal(self) -> (MarketEvent, Option<u64>) {
186        (self, None)
187    }
188}
189
190impl EventSourceItem for SequencedMarketEvent {
191    fn into_event_and_ordinal(self) -> (MarketEvent, Option<u64>) {
192        (self.event, Some(self.source_row_ordinal))
193    }
194}
195
196/// A contiguous batch of events with the same timestamp.
197#[derive(Debug, Clone, Serialize, Deserialize)]
198pub struct TimestampBatch {
199    pub ts: NaiveDateTime,
200    pub events: Vec<FeedEvent>,
201}
202
203impl TimestampBatch {
204    pub fn len(&self) -> usize {
205        self.events.len()
206    }
207
208    pub fn is_empty(&self) -> bool {
209        self.events.is_empty()
210    }
211}
212
213/// Sequential source of market events for backtesting.
214pub trait DataFeed {
215    /// Return the next event, or `None` when the feed is exhausted.
216    fn next_event(&mut self) -> Option<MarketEvent>;
217
218    /// Peek at the next event without consuming it.
219    fn peek(&self) -> Option<&MarketEvent>;
220
221    /// Return all contiguous events at the next timestamp.
222    ///
223    /// Implementations with native metadata should override this method.
224    /// The default preserves compatibility for custom feeds and assigns primary role metadata in their existing event order.
225    fn next_batch(&mut self) -> Option<TimestampBatch> {
226        let ts = self.peek()?.ts();
227        let mut events = Vec::new();
228        let mut row_sequence = 0;
229
230        while self.peek().is_some_and(|event| event.ts() == ts) {
231            let event = self.next_event()?;
232            events.push(FeedEvent::new(
233                event,
234                EventMetadata::new(SeriesRoles::PRIMARY, 0, row_sequence),
235            ));
236            row_sequence += 1;
237        }
238
239        Some(TimestampBatch { ts, events })
240    }
241
242    /// Total number of events when known without consuming the feed.
243    fn total_events(&self) -> Option<usize> {
244        None
245    }
246}
247
248/// Fallible source of complete timestamp batches.
249pub trait FallibleBatchFeed {
250    type Error;
251
252    /// Return every event at the next timestamp, or `None` at end of input.
253    fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error>;
254}
255
256impl<F> FallibleBatchFeed for Box<F>
257where
258    F: FallibleBatchFeed + ?Sized,
259{
260    type Error = F::Error;
261
262    fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
263        (**self).next_batch()
264    }
265}
266
267/// Errors produced while grouping a fallible event source into timestamp batches.
268#[derive(Debug, Error)]
269pub enum EventBatchFeedError<E> {
270    #[error("event source failed: {0}")]
271    Source(E),
272    #[error("event source moved backwards from {previous} to {current}")]
273    NonMonotonic {
274        previous: NaiveDateTime,
275        current: NaiveDateTime,
276    },
277    #[error("event row sequence is exhausted")]
278    RowSequenceExhausted,
279}
280
281/// Groups a fallible event cursor into complete timestamp batches.
282pub struct EventBatchFeed<F, I = MarketEvent> {
283    next_event: F,
284    pending: Option<FeedEvent>,
285    roles: SeriesRoles,
286    series_rank: u32,
287    next_row_sequence: u64,
288    last_source_ts: Option<NaiveDateTime>,
289    exhausted: bool,
290    item: PhantomData<fn() -> I>,
291}
292
293impl<F, I> EventBatchFeed<F, I> {
294    pub fn new<E>(next_event: F, roles: SeriesRoles, series_rank: u32) -> Self
295    where
296        F: FnMut() -> Result<Option<I>, E>,
297        I: EventSourceItem,
298    {
299        Self {
300            next_event,
301            pending: None,
302            roles,
303            series_rank,
304            next_row_sequence: 0,
305            last_source_ts: None,
306            exhausted: false,
307            item: PhantomData,
308        }
309    }
310}
311
312impl<F, I, E> EventBatchFeed<F, I>
313where
314    F: FnMut() -> Result<Option<I>, E>,
315    I: EventSourceItem,
316{
317    /// Return the next complete timestamp batch.
318    pub fn next_timestamp_batch(
319        &mut self,
320    ) -> Result<Option<TimestampBatch>, EventBatchFeedError<E>> {
321        let first = match self.pending.take() {
322            Some(event) => event,
323            None => match self.pull_event()? {
324                Some(event) => event,
325                None => return Ok(None),
326            },
327        };
328        let ts = first.event.ts();
329        let mut events = vec![first];
330
331        loop {
332            match self.pull_event()? {
333                Some(event) if event.event.ts() == ts => events.push(event),
334                Some(event) => {
335                    self.pending = Some(event);
336                    break;
337                }
338                None => break,
339            }
340        }
341
342        Ok(Some(TimestampBatch { ts, events }))
343    }
344
345    fn pull_event(&mut self) -> Result<Option<FeedEvent>, EventBatchFeedError<E>> {
346        if self.exhausted {
347            return Ok(None);
348        }
349        let item = match (self.next_event)() {
350            Ok(Some(item)) => item,
351            Ok(None) => {
352                self.exhausted = true;
353                return Ok(None);
354            }
355            Err(error) => {
356                self.exhausted = true;
357                return Err(EventBatchFeedError::Source(error));
358            }
359        };
360        let (event, source_row_ordinal) = item.into_event_and_ordinal();
361        let current = event.ts();
362        if let Some(previous) = self.last_source_ts
363            && current < previous
364        {
365            self.exhausted = true;
366            return Err(EventBatchFeedError::NonMonotonic { previous, current });
367        }
368        let row_sequence = match source_row_ordinal {
369            Some(source_row_ordinal) => source_row_ordinal,
370            None => {
371                let row_sequence = self.next_row_sequence;
372                self.next_row_sequence = self
373                    .next_row_sequence
374                    .checked_add(1)
375                    .ok_or(EventBatchFeedError::RowSequenceExhausted)?;
376                row_sequence
377            }
378        };
379        self.last_source_ts = Some(current);
380        Ok(Some(FeedEvent::new(
381            event,
382            EventMetadata::new(self.roles, self.series_rank, row_sequence),
383        )))
384    }
385}
386
387impl<F, I, E> FallibleBatchFeed for EventBatchFeed<F, I>
388where
389    F: FnMut() -> Result<Option<I>, E>,
390    I: EventSourceItem,
391{
392    type Error = EventBatchFeedError<E>;
393
394    fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
395        self.next_timestamp_batch()
396    }
397}
398
399/// Errors produced while deterministically merging timestamp batch feeds.
400#[derive(Debug, Error)]
401pub enum KWayMergeError<E> {
402    #[error("series {series_rank} failed: {error}")]
403    Source { series_rank: u32, error: E },
404    #[error("series {series_rank} returned an empty timestamp batch at {ts}")]
405    EmptyBatch { series_rank: u32, ts: NaiveDateTime },
406    #[error("series {series_rank} batch at {batch_ts} contains an event at {event_ts}")]
407    TimestampMismatch {
408        series_rank: u32,
409        batch_ts: NaiveDateTime,
410        event_ts: NaiveDateTime,
411    },
412    #[error("series {series_rank} moved backwards from {previous} to {current}")]
413    NonMonotonic {
414        series_rank: u32,
415        previous: NaiveDateTime,
416        current: NaiveDateTime,
417    },
418    #[error("duplicate merge ordering key ({ts}, {series_rank}, {row_sequence})")]
419    DuplicateOrderingKey {
420        ts: NaiveDateTime,
421        series_rank: u32,
422        row_sequence: u64,
423    },
424}
425
426/// Streaming deterministic merge over multiple fallible timestamp batch feeds.
427pub struct KWayMergeFeed<F>
428where
429    F: FallibleBatchFeed,
430{
431    feeds: Vec<F>,
432    heads: Vec<Option<TimestampBatch>>,
433    head_heap: BinaryHeap<Reverse<(NaiveDateTime, usize)>>,
434    exhausted: Vec<bool>,
435    last_batch_ts: Vec<Option<NaiveDateTime>>,
436    initialized: bool,
437}
438
439impl<F> KWayMergeFeed<F>
440where
441    F: FallibleBatchFeed,
442{
443    /// Create a merge using input position as `series_rank`.
444    pub fn new(feeds: Vec<F>) -> Self {
445        assert!(
446            feeds.len() <= u32::MAX as usize,
447            "KWayMergeFeed supports at most u32::MAX series"
448        );
449        let series_count = feeds.len();
450        Self {
451            feeds,
452            heads: (0..series_count).map(|_| None).collect(),
453            head_heap: BinaryHeap::with_capacity(series_count),
454            exhausted: vec![false; series_count],
455            last_batch_ts: vec![None; series_count],
456            initialized: false,
457        }
458    }
459
460    pub fn series_count(&self) -> usize {
461        self.feeds.len()
462    }
463
464    pub fn is_empty(&self) -> bool {
465        self.feeds.is_empty()
466    }
467
468    /// Return the next complete merged timestamp batch.
469    pub fn next_timestamp_batch(
470        &mut self,
471    ) -> Result<Option<TimestampBatch>, KWayMergeError<F::Error>> {
472        if !self.initialized {
473            for index in 0..self.feeds.len() {
474                self.fill_head(index)?;
475            }
476            self.initialized = true;
477        }
478        let Some(Reverse((ts, _))) = self.head_heap.peek().copied() else {
479            return Ok(None);
480        };
481
482        let mut events = Vec::new();
483        while self
484            .head_heap
485            .peek()
486            .is_some_and(|Reverse((head_ts, _))| *head_ts == ts)
487        {
488            let Reverse((_, index)) = self.head_heap.pop().expect("matching heap head exists");
489            let batch = self.heads[index]
490                .take()
491                .expect("heap head retains its batch");
492            let series_rank = index as u32;
493            events.extend(batch.events.into_iter().map(|mut event| {
494                event.metadata.series_rank = series_rank;
495                event
496            }));
497            self.fill_head(index)?;
498        }
499
500        events.sort_by_key(FeedEvent::ordering_key);
501        for duplicate in events.windows(2) {
502            if duplicate[0].ordering_key() == duplicate[1].ordering_key() {
503                let (_, series_rank, row_sequence) = duplicate[0].ordering_key();
504                return Err(KWayMergeError::DuplicateOrderingKey {
505                    ts,
506                    series_rank,
507                    row_sequence,
508                });
509            }
510        }
511        Ok(Some(TimestampBatch { ts, events }))
512    }
513
514    fn fill_head(&mut self, index: usize) -> Result<(), KWayMergeError<F::Error>> {
515        if self.heads[index].is_some() || self.exhausted[index] {
516            return Ok(());
517        }
518        let series_rank = index as u32;
519        let Some(batch) = self.feeds[index]
520            .next_batch()
521            .map_err(|error| KWayMergeError::Source { series_rank, error })?
522        else {
523            self.exhausted[index] = true;
524            return Ok(());
525        };
526        if batch.is_empty() {
527            return Err(KWayMergeError::EmptyBatch {
528                series_rank,
529                ts: batch.ts,
530            });
531        }
532        for event in &batch.events {
533            if event.event.ts() != batch.ts {
534                return Err(KWayMergeError::TimestampMismatch {
535                    series_rank,
536                    batch_ts: batch.ts,
537                    event_ts: event.event.ts(),
538                });
539            }
540        }
541        if let Some(previous) = self.last_batch_ts[index]
542            && batch.ts < previous
543        {
544            return Err(KWayMergeError::NonMonotonic {
545                series_rank,
546                previous,
547                current: batch.ts,
548            });
549        }
550        self.last_batch_ts[index] = Some(batch.ts);
551        self.head_heap.push(Reverse((batch.ts, index)));
552        self.heads[index] = Some(batch);
553        Ok(())
554    }
555}
556
557impl<F> FallibleBatchFeed for KWayMergeFeed<F>
558where
559    F: FallibleBatchFeed,
560{
561    type Error = KWayMergeError<F::Error>;
562
563    fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
564        self.next_timestamp_batch()
565    }
566}
567
568/// In-memory data feed backed by metadata-bearing market events.
569///
570/// [`Self::new`] retains the existing contract and does not sort its input.
571/// [`Self::from_feed_events`] sorts by `(timestamp, series_rank, row_sequence)`.
572#[derive(Debug, Clone)]
573pub struct VecFeed {
574    events: Vec<FeedEvent>,
575    index: usize,
576}
577
578impl VecFeed {
579    /// Create a new feed from a pre-sorted vector of primary events.
580    pub fn new(events: Vec<MarketEvent>) -> Self {
581        let events = events
582            .into_iter()
583            .enumerate()
584            .map(|(row_sequence, event)| {
585                FeedEvent::new(
586                    event,
587                    EventMetadata::new(SeriesRoles::PRIMARY, 0, row_sequence as u64),
588                )
589            })
590            .collect();
591        Self { events, index: 0 }
592    }
593
594    /// Create a deterministically sorted feed from metadata-bearing events.
595    pub fn from_feed_events(mut events: Vec<FeedEvent>) -> Self {
596        events.sort_by_key(FeedEvent::ordering_key);
597        Self { events, index: 0 }
598    }
599
600    /// Return the next event together with its ordering metadata.
601    pub fn next_feed_event(&mut self) -> Option<FeedEvent> {
602        let event = self.events.get(self.index)?.clone();
603        self.index += 1;
604        Some(event)
605    }
606
607    /// Peek at the metadata for the next event without consuming it.
608    pub fn peek_metadata(&self) -> Option<&EventMetadata> {
609        self.events.get(self.index).map(|event| &event.metadata)
610    }
611
612    /// Return all remaining events at the next timestamp.
613    pub fn next_timestamp_batch(&mut self) -> Option<TimestampBatch> {
614        let ts = self.events.get(self.index)?.event.ts();
615        let start = self.index;
616        while self
617            .events
618            .get(self.index)
619            .is_some_and(|event| event.event.ts() == ts)
620        {
621            self.index += 1;
622        }
623
624        Some(TimestampBatch {
625            ts,
626            events: self.events[start..self.index].to_vec(),
627        })
628    }
629
630    /// Number of events remaining.
631    pub fn remaining(&self) -> usize {
632        self.events.len().saturating_sub(self.index)
633    }
634
635    /// Total number of events in the feed (consumed + remaining).
636    pub fn total(&self) -> usize {
637        self.events.len()
638    }
639
640    /// Reset the feed to the beginning.
641    pub fn reset(&mut self) {
642        self.index = 0;
643    }
644}
645
646impl DataFeed for VecFeed {
647    fn next_event(&mut self) -> Option<MarketEvent> {
648        self.next_feed_event().map(|event| event.event)
649    }
650
651    fn peek(&self) -> Option<&MarketEvent> {
652        self.events.get(self.index).map(|event| &event.event)
653    }
654
655    fn next_batch(&mut self) -> Option<TimestampBatch> {
656        self.next_timestamp_batch()
657    }
658
659    fn total_events(&self) -> Option<usize> {
660        Some(self.events.len())
661    }
662}
663
664/// Convert ticks into a primary-series feed.
665///
666/// Ticks without a valid executable bid/ask quote are silently skipped.
667pub fn ticks_to_feed(ticks: Vec<data_preprocess::Tick>) -> VecFeed {
668    ticks_to_feed_with_metadata(ticks, SeriesRoles::PRIMARY, 0)
669}
670
671/// Convert ticks into a deterministically ranked feed.
672pub fn ticks_to_feed_with_metadata(
673    ticks: Vec<data_preprocess::Tick>,
674    roles: SeriesRoles,
675    series_rank: u32,
676) -> VecFeed {
677    let events = ticks
678        .into_iter()
679        .enumerate()
680        .filter_map(|(row_sequence, tick)| {
681            let bid = tick.bid?;
682            let ask = tick.ask?;
683            let event = MarketEvent::Tick {
684                symbol: tick.symbol,
685                ts: tick.ts,
686                bid,
687                ask,
688            };
689            event.to_valid_quote().map(|_| {
690                FeedEvent::new(
691                    event,
692                    EventMetadata::new(roles, series_rank, row_sequence as u64),
693                )
694            })
695        })
696        .collect();
697    VecFeed::from_feed_events(events)
698}
699
700/// Convert bars into a primary-series feed.
701pub fn bars_to_feed(bars: Vec<data_preprocess::Bar>) -> VecFeed {
702    bars_to_feed_with_metadata(bars, SeriesRoles::PRIMARY, 0)
703}
704
705/// Convert bars into a deterministically ranked feed.
706pub fn bars_to_feed_with_metadata(
707    bars: Vec<data_preprocess::Bar>,
708    roles: SeriesRoles,
709    series_rank: u32,
710) -> VecFeed {
711    let events = bars
712        .into_iter()
713        .enumerate()
714        .map(|(row_sequence, bar)| {
715            let event = MarketEvent::Bar {
716                symbol: bar.symbol,
717                ts: bar.ts,
718                open: bar.open,
719                high: bar.high,
720                low: bar.low,
721                close: bar.close,
722                volume: bar.volume,
723            };
724            FeedEvent::new(
725                event,
726                EventMetadata::new(roles, series_rank, row_sequence as u64),
727            )
728        })
729        .collect();
730    VecFeed::from_feed_events(events)
731}
732
733/// Merge feeds using input position as `series_rank`.
734///
735/// Events are ordered by `(timestamp, series_rank, row_sequence)`.
736/// Role metadata is preserved, so a single event can serve both primary and conversion roles.
737pub fn merge_feeds(feeds: Vec<VecFeed>) -> VecFeed {
738    let mut all_events = Vec::new();
739    for (series_rank, feed) in feeds.into_iter().enumerate() {
740        let series_rank = series_rank.min(u32::MAX as usize) as u32;
741        all_events.extend(feed.events.into_iter().map(|mut event| {
742            event.metadata.series_rank = series_rank;
743            event
744        }));
745    }
746    VecFeed::from_feed_events(all_events)
747}
748
749// ─── Tests ──────────────────────────────────────────────────────────────────
750
751#[cfg(test)]
752mod tests {
753    use std::collections::VecDeque;
754
755    use super::*;
756    use chrono::NaiveDate;
757
758    fn ts(h: u32, m: u32, s: u32) -> NaiveDateTime {
759        NaiveDate::from_ymd_opt(2026, 1, 1)
760            .unwrap()
761            .and_hms_opt(h, m, s)
762            .unwrap()
763    }
764
765    fn sample_events() -> Vec<MarketEvent> {
766        vec![
767            MarketEvent::Tick {
768                symbol: "EURUSD".into(),
769                ts: ts(10, 0, 0),
770                bid: 1.0848,
771                ask: 1.0850,
772            },
773            MarketEvent::Tick {
774                symbol: "EURUSD".into(),
775                ts: ts(10, 0, 1),
776                bid: 1.0849,
777                ask: 1.0851,
778            },
779            MarketEvent::Tick {
780                symbol: "EURUSD".into(),
781                ts: ts(10, 0, 2),
782                bid: 1.0847,
783                ask: 1.0849,
784            },
785        ]
786    }
787
788    #[test]
789    fn vec_feed_iterates() {
790        let mut feed = VecFeed::new(sample_events());
791        assert_eq!(feed.total(), 3);
792        assert_eq!(feed.remaining(), 3);
793
794        let e1 = feed.next_event().unwrap();
795        assert_eq!(e1.symbol(), "EURUSD");
796        assert_eq!(feed.remaining(), 2);
797
798        let _ = feed.next_event().unwrap();
799        let _ = feed.next_event().unwrap();
800        assert!(feed.next_event().is_none());
801        assert_eq!(feed.remaining(), 0);
802    }
803
804    #[test]
805    fn vec_feed_peek() {
806        let feed = VecFeed::new(sample_events());
807        let peeked = feed.peek().unwrap();
808        assert_eq!(peeked.ts(), ts(10, 0, 0));
809    }
810
811    #[test]
812    fn vec_feed_reset() {
813        let mut feed = VecFeed::new(sample_events());
814        let _ = feed.next_event();
815        let _ = feed.next_event();
816        feed.reset();
817        assert_eq!(feed.remaining(), 3);
818    }
819
820    #[test]
821    fn to_quote_tick() {
822        let event = MarketEvent::Tick {
823            symbol: "EURUSD".into(),
824            ts: ts(10, 0, 0),
825            bid: 1.0848,
826            ask: 1.0850,
827        };
828        let q = event.to_quote();
829        assert_eq!(q.symbol, "EURUSD");
830        assert!((q.bid - 1.0848).abs() < f64::EPSILON);
831        assert!((q.ask - 1.0850).abs() < f64::EPSILON);
832    }
833
834    #[test]
835    fn to_quote_bar() {
836        let event = MarketEvent::Bar {
837            symbol: "EURUSD".into(),
838            ts: ts(10, 0, 0),
839            open: 1.0840,
840            high: 1.0860,
841            low: 1.0830,
842            close: 1.0855,
843            volume: 1000,
844        };
845        let q = event.to_quote();
846        // Bar uses close for both bid and ask
847        assert!((q.bid - 1.0855).abs() < f64::EPSILON);
848        assert!((q.ask - 1.0855).abs() < f64::EPSILON);
849    }
850
851    #[test]
852    fn merge_two_feeds() {
853        let feed_a = VecFeed::new(vec![
854            MarketEvent::Tick {
855                symbol: "EURUSD".into(),
856                ts: ts(10, 0, 0),
857                bid: 1.08,
858                ask: 1.09,
859            },
860            MarketEvent::Tick {
861                symbol: "EURUSD".into(),
862                ts: ts(10, 0, 2),
863                bid: 1.08,
864                ask: 1.09,
865            },
866        ]);
867        let feed_b = VecFeed::new(vec![MarketEvent::Tick {
868            symbol: "XAUUSD".into(),
869            ts: ts(10, 0, 1),
870            bid: 2000.0,
871            ask: 2001.0,
872        }]);
873
874        let mut merged = merge_feeds(vec![feed_a, feed_b]);
875        assert_eq!(merged.total(), 3);
876
877        let e1 = merged.next_event().unwrap();
878        assert_eq!(e1.ts(), ts(10, 0, 0));
879        assert_eq!(e1.symbol(), "EURUSD");
880
881        let e2 = merged.next_event().unwrap();
882        assert_eq!(e2.ts(), ts(10, 0, 1));
883        assert_eq!(e2.symbol(), "XAUUSD");
884
885        let e3 = merged.next_event().unwrap();
886        assert_eq!(e3.ts(), ts(10, 0, 2));
887        assert_eq!(e3.symbol(), "EURUSD");
888    }
889
890    #[test]
891    fn ticks_to_feed_skips_missing_prices() {
892        let ticks = vec![
893            data_preprocess::Tick {
894                exchange: "test".into(),
895                symbol: "EURUSD".into(),
896                ts: ts(10, 0, 0),
897                bid: Some(1.08),
898                ask: Some(1.09),
899                last: None,
900                volume: None,
901                flags: None,
902            },
903            data_preprocess::Tick {
904                exchange: "test".into(),
905                symbol: "EURUSD".into(),
906                ts: ts(10, 0, 1),
907                bid: Some(1.08),
908                ask: None, // missing ask
909                last: None,
910                volume: None,
911                flags: None,
912            },
913            data_preprocess::Tick {
914                exchange: "test".into(),
915                symbol: "EURUSD".into(),
916                ts: ts(10, 0, 2),
917                bid: None, // missing bid
918                ask: Some(1.09),
919                last: None,
920                volume: None,
921                flags: None,
922            },
923        ];
924
925        let feed = ticks_to_feed(ticks);
926        assert_eq!(feed.total(), 1); // only the first tick survives
927    }
928
929    #[test]
930    fn validated_quote_contract_rejects_nonfinite_nonpositive_and_crossed_ticks() {
931        for (bid, ask) in [
932            (f64::NAN, 1.0),
933            (1.0, f64::INFINITY),
934            (0.0, 1.0),
935            (2.0, 1.0),
936        ] {
937            let event = MarketEvent::Tick {
938                symbol: "EURUSD".into(),
939                ts: ts(10, 0, 0),
940                bid,
941                ask,
942            };
943            assert!(event.to_valid_quote().is_none());
944        }
945
946        let ticks = vec![
947            data_preprocess::Tick {
948                exchange: "test".into(),
949                symbol: "EURUSD".into(),
950                ts: ts(10, 0, 0),
951                bid: Some(1.1),
952                ask: Some(1.0),
953                last: None,
954                volume: None,
955                flags: None,
956            },
957            data_preprocess::Tick {
958                exchange: "test".into(),
959                symbol: "EURUSD".into(),
960                ts: ts(10, 0, 1),
961                bid: Some(1.0),
962                ask: Some(1.1),
963                last: None,
964                volume: None,
965                flags: None,
966            },
967        ];
968        assert_eq!(ticks_to_feed(ticks).total(), 1);
969    }
970
971    #[test]
972    fn bars_to_feed_converts_all() {
973        let bars = vec![
974            data_preprocess::Bar {
975                exchange: "test".into(),
976                symbol: "EURUSD".into(),
977                timeframe: data_preprocess::Timeframe::M5,
978                ts: ts(10, 0, 0),
979                open: 1.0840,
980                high: 1.0860,
981                low: 1.0830,
982                close: 1.0855,
983                tick_vol: 100,
984                volume: 1000,
985                spread: 2,
986            },
987            data_preprocess::Bar {
988                exchange: "test".into(),
989                symbol: "EURUSD".into(),
990                timeframe: data_preprocess::Timeframe::M5,
991                ts: ts(10, 5, 0),
992                open: 1.0855,
993                high: 1.0870,
994                low: 1.0845,
995                close: 1.0865,
996                tick_vol: 120,
997                volume: 1200,
998                spread: 2,
999            },
1000        ];
1001
1002        let feed = bars_to_feed(bars);
1003        assert_eq!(feed.total(), 2);
1004    }
1005
1006    #[test]
1007    fn equal_timestamp_batch_uses_deterministic_metadata_order() {
1008        let batch_ts = ts(10, 0, 0);
1009        let later_ts = ts(10, 0, 1);
1010        let events = vec![
1011            FeedEvent::new(
1012                MarketEvent::Tick {
1013                    symbol: "CONVERSION".into(),
1014                    ts: batch_ts,
1015                    bid: 2.0,
1016                    ask: 2.1,
1017                },
1018                EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
1019            ),
1020            FeedEvent::new(
1021                MarketEvent::Bar {
1022                    symbol: "PRIMARY".into(),
1023                    ts: batch_ts,
1024                    open: 1.0,
1025                    high: 1.2,
1026                    low: 0.9,
1027                    close: 1.1,
1028                    volume: 10,
1029                },
1030                EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
1031            ),
1032            FeedEvent::new(
1033                MarketEvent::Tick {
1034                    symbol: "PRIMARY".into(),
1035                    ts: batch_ts,
1036                    bid: 1.0,
1037                    ask: 1.1,
1038                },
1039                EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
1040            ),
1041            FeedEvent::new(
1042                MarketEvent::Tick {
1043                    symbol: "PRIMARY".into(),
1044                    ts: later_ts,
1045                    bid: 1.1,
1046                    ask: 1.2,
1047                },
1048                EventMetadata::new(SeriesRoles::PRIMARY, 0, 2),
1049            ),
1050        ];
1051        let mut feed = VecFeed::from_feed_events(events);
1052
1053        let batch = feed
1054            .next_timestamp_batch()
1055            .expect("equal timestamp batch should be available");
1056        assert_eq!(batch.ts, batch_ts);
1057        assert_eq!(batch.len(), 3);
1058        let keys: Vec<_> = batch.events.iter().map(FeedEvent::ordering_key).collect();
1059        assert_eq!(
1060            keys,
1061            vec![(batch_ts, 0, 0), (batch_ts, 0, 1), (batch_ts, 1, 0)]
1062        );
1063        assert!(matches!(batch.events[1].event, MarketEvent::Bar { .. }));
1064        assert!(matches!(batch.events[2].event, MarketEvent::Tick { .. }));
1065
1066        let later = DataFeed::next_batch(&mut feed).expect("later batch should remain");
1067        assert_eq!(later.ts, later_ts);
1068        assert_eq!(later.len(), 1);
1069        assert!(DataFeed::next_batch(&mut feed).is_none());
1070    }
1071
1072    #[test]
1073    fn one_tick_can_serve_primary_and_conversion_roles() {
1074        let ticks = vec![data_preprocess::Tick {
1075            exchange: "test".into(),
1076            symbol: "EURUSD".into(),
1077            ts: ts(10, 0, 0),
1078            bid: Some(1.08),
1079            ask: Some(1.09),
1080            last: None,
1081            volume: None,
1082            flags: None,
1083        }];
1084        let mut feed = ticks_to_feed_with_metadata(ticks, SeriesRoles::PRIMARY_AND_CONVERSION, 0);
1085
1086        let batch = feed.next_timestamp_batch().unwrap();
1087        assert_eq!(batch.len(), 1);
1088        assert_eq!(
1089            batch.events[0].metadata.roles,
1090            SeriesRoles::PRIMARY_AND_CONVERSION
1091        );
1092        assert_eq!(batch.events[0].metadata.row_sequence, 0);
1093    }
1094
1095    #[test]
1096    fn merge_feeds_ranks_series_and_batches_equal_timestamps_stably() {
1097        let event_ts = ts(10, 0, 0);
1098        let primary = VecFeed::new(vec![MarketEvent::Bar {
1099            symbol: "PRIMARY".into(),
1100            ts: event_ts,
1101            open: 1.0,
1102            high: 1.2,
1103            low: 0.9,
1104            close: 1.1,
1105            volume: 10,
1106        }]);
1107        let conversion = VecFeed::from_feed_events(vec![FeedEvent::new(
1108            MarketEvent::Tick {
1109                symbol: "CONVERSION".into(),
1110                ts: event_ts,
1111                bid: 2.0,
1112                ask: 2.1,
1113            },
1114            EventMetadata::new(SeriesRoles::CONVERSION, 0, 0),
1115        )]);
1116
1117        let mut merged = merge_feeds(vec![primary, conversion]);
1118        let batch = merged.next_timestamp_batch().unwrap();
1119        assert_eq!(batch.len(), 2);
1120        assert_eq!(batch.events[0].metadata.series_rank, 0);
1121        assert_eq!(batch.events[0].metadata.roles, SeriesRoles::PRIMARY);
1122        assert_eq!(batch.events[1].metadata.series_rank, 1);
1123        assert_eq!(batch.events[1].metadata.roles, SeriesRoles::CONVERSION);
1124    }
1125
1126    fn ranked_tick(
1127        symbol: &str,
1128        event_ts: NaiveDateTime,
1129        roles: SeriesRoles,
1130        row_sequence: u64,
1131    ) -> FeedEvent {
1132        FeedEvent::new(
1133            MarketEvent::Tick {
1134                symbol: symbol.into(),
1135                ts: event_ts,
1136                bid: 1.0,
1137                ask: 1.1,
1138            },
1139            EventMetadata::new(roles, 99, row_sequence),
1140        )
1141    }
1142
1143    struct ScriptedFeed {
1144        batches: VecDeque<Result<TimestampBatch, &'static str>>,
1145    }
1146
1147    impl ScriptedFeed {
1148        fn new(batches: Vec<Result<TimestampBatch, &'static str>>) -> Self {
1149            Self {
1150                batches: batches.into(),
1151            }
1152        }
1153    }
1154
1155    impl FallibleBatchFeed for ScriptedFeed {
1156        type Error = &'static str;
1157
1158        fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
1159            self.batches.pop_front().transpose()
1160        }
1161    }
1162
1163    #[test]
1164    fn event_batch_feed_groups_complete_timestamps_and_supports_shared_roles() {
1165        let first_ts = ts(10, 0, 0);
1166        let second_ts = ts(10, 0, 1);
1167        let mut source = vec![
1168            MarketEvent::Tick {
1169                symbol: "EURUSD".into(),
1170                ts: first_ts,
1171                bid: 1.0,
1172                ask: 1.1,
1173            },
1174            MarketEvent::Bar {
1175                symbol: "EURUSD".into(),
1176                ts: first_ts,
1177                open: 1.0,
1178                high: 1.2,
1179                low: 0.9,
1180                close: 1.1,
1181                volume: 10,
1182            },
1183            MarketEvent::Tick {
1184                symbol: "EURUSD".into(),
1185                ts: second_ts,
1186                bid: 1.1,
1187                ask: 1.2,
1188            },
1189        ]
1190        .into_iter();
1191        let mut feed = EventBatchFeed::new(
1192            move || Ok::<_, &'static str>(source.next()),
1193            SeriesRoles::PRIMARY_AND_CONVERSION,
1194            4,
1195        );
1196
1197        let first = feed.next_timestamp_batch().unwrap().unwrap();
1198        assert_eq!(first.ts, first_ts);
1199        assert_eq!(first.len(), 2);
1200        assert!(
1201            first
1202                .events
1203                .iter()
1204                .all(|event| event.metadata.roles == SeriesRoles::PRIMARY_AND_CONVERSION)
1205        );
1206        assert_eq!(first.events[0].metadata.row_sequence, 0);
1207        assert_eq!(first.events[1].metadata.row_sequence, 1);
1208
1209        let second = feed.next_timestamp_batch().unwrap().unwrap();
1210        assert_eq!(second.ts, second_ts);
1211        assert_eq!(second.events[0].metadata.row_sequence, 2);
1212        assert!(feed.next_timestamp_batch().unwrap().is_none());
1213    }
1214
1215    #[test]
1216    fn event_batch_feed_preserves_explicit_physical_source_ordinals() {
1217        let mut source = vec![
1218            SequencedMarketEvent::new(
1219                MarketEvent::Tick {
1220                    symbol: "EURUSD".into(),
1221                    ts: ts(10, 0, 0),
1222                    bid: 1.0,
1223                    ask: 1.1,
1224                },
1225                4,
1226            ),
1227            SequencedMarketEvent::new(
1228                MarketEvent::Tick {
1229                    symbol: "EURUSD".into(),
1230                    ts: ts(10, 0, 0),
1231                    bid: 1.1,
1232                    ask: 1.2,
1233                },
1234                9,
1235            ),
1236        ]
1237        .into_iter();
1238        let mut feed = EventBatchFeed::new(
1239            move || Ok::<_, &'static str>(source.next()),
1240            SeriesRoles::PRIMARY,
1241            2,
1242        );
1243
1244        let batch = feed.next_timestamp_batch().unwrap().unwrap();
1245        assert_eq!(
1246            batch
1247                .events
1248                .iter()
1249                .map(|event| event.metadata.row_sequence)
1250                .collect::<Vec<_>>(),
1251            vec![4, 9]
1252        );
1253        assert_eq!(batch.events[0].metadata.series_rank, 2);
1254    }
1255
1256    #[test]
1257    fn event_batch_feed_rejects_a_non_monotonic_event_source() {
1258        let mut source = vec![
1259            MarketEvent::Tick {
1260                symbol: "EURUSD".into(),
1261                ts: ts(10, 0, 1),
1262                bid: 1.0,
1263                ask: 1.1,
1264            },
1265            MarketEvent::Tick {
1266                symbol: "EURUSD".into(),
1267                ts: ts(10, 0, 0),
1268                bid: 1.0,
1269                ask: 1.1,
1270            },
1271        ]
1272        .into_iter();
1273        let mut feed = EventBatchFeed::new(
1274            move || Ok::<_, &'static str>(source.next()),
1275            SeriesRoles::PRIMARY,
1276            0,
1277        );
1278
1279        assert!(matches!(
1280            feed.next_timestamp_batch(),
1281            Err(EventBatchFeedError::NonMonotonic { previous, current })
1282                if previous == ts(10, 0, 1) && current == ts(10, 0, 0)
1283        ));
1284    }
1285
1286    #[test]
1287    fn k_way_merge_drains_all_equal_timestamp_batches_and_orders_by_full_key() {
1288        let first_ts = ts(10, 0, 0);
1289        let second_ts = ts(10, 0, 1);
1290        let first = ScriptedFeed::new(vec![
1291            Ok(TimestampBatch {
1292                ts: first_ts,
1293                events: vec![ranked_tick(
1294                    "SHARED",
1295                    first_ts,
1296                    SeriesRoles::PRIMARY_AND_CONVERSION,
1297                    0,
1298                )],
1299            }),
1300            Ok(TimestampBatch {
1301                ts: first_ts,
1302                events: vec![ranked_tick("PRIMARY", first_ts, SeriesRoles::PRIMARY, 1)],
1303            }),
1304            Ok(TimestampBatch {
1305                ts: second_ts,
1306                events: vec![ranked_tick("PRIMARY", second_ts, SeriesRoles::PRIMARY, 2)],
1307            }),
1308        ]);
1309        let second = ScriptedFeed::new(vec![Ok(TimestampBatch {
1310            ts: first_ts,
1311            events: vec![ranked_tick(
1312                "CONVERSION",
1313                first_ts,
1314                SeriesRoles::CONVERSION,
1315                0,
1316            )],
1317        })]);
1318        let mut merged = KWayMergeFeed::new(vec![first, second]);
1319
1320        let batch = merged.next_timestamp_batch().unwrap().unwrap();
1321        assert_eq!(batch.ts, first_ts);
1322        assert_eq!(batch.len(), 3);
1323        assert_eq!(
1324            batch
1325                .events
1326                .iter()
1327                .map(FeedEvent::ordering_key)
1328                .collect::<Vec<_>>(),
1329            vec![(first_ts, 0, 0), (first_ts, 0, 1), (first_ts, 1, 0)]
1330        );
1331        assert_eq!(
1332            batch.events[0].metadata.roles,
1333            SeriesRoles::PRIMARY_AND_CONVERSION
1334        );
1335
1336        let later = merged.next_timestamp_batch().unwrap().unwrap();
1337        assert_eq!(later.ts, second_ts);
1338        assert_eq!(later.len(), 1);
1339        assert!(merged.next_timestamp_batch().unwrap().is_none());
1340    }
1341
1342    #[test]
1343    fn k_way_merge_propagates_errors_while_draining_a_complete_timestamp() {
1344        let event_ts = ts(10, 0, 0);
1345        let first = ScriptedFeed::new(vec![
1346            Ok(TimestampBatch {
1347                ts: event_ts,
1348                events: vec![ranked_tick("PRIMARY", event_ts, SeriesRoles::PRIMARY, 0)],
1349            }),
1350            Err("same timestamp continuation failed"),
1351        ]);
1352        let second = ScriptedFeed::new(vec![Ok(TimestampBatch {
1353            ts: event_ts,
1354            events: vec![ranked_tick(
1355                "CONVERSION",
1356                event_ts,
1357                SeriesRoles::CONVERSION,
1358                0,
1359            )],
1360        })]);
1361        let mut merged = KWayMergeFeed::new(vec![first, second]);
1362
1363        assert!(matches!(
1364            merged.next_timestamp_batch(),
1365            Err(KWayMergeError::Source {
1366                series_rank: 0,
1367                error: "same timestamp continuation failed"
1368            })
1369        ));
1370    }
1371
1372    #[test]
1373    fn k_way_merge_propagates_ranked_source_failures() {
1374        let mut merged = KWayMergeFeed::new(vec![ScriptedFeed::new(vec![Err("read failed")])]);
1375        assert!(matches!(
1376            merged.next_timestamp_batch(),
1377            Err(KWayMergeError::Source {
1378                series_rank: 0,
1379                error: "read failed"
1380            })
1381        ));
1382    }
1383}