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