1use 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#[derive(Debug, Clone, Serialize, Deserialize)]
25pub enum MarketEvent {
26 Tick {
28 symbol: String,
29 ts: NaiveDateTime,
30 bid: f64,
31 ask: f64,
32 },
33 Bar {
35 symbol: String,
36 ts: NaiveDateTime,
37 open: f64,
38 high: f64,
39 low: f64,
40 close: f64,
41 volume: i64,
42 #[serde(default)]
46 spread: Option<f64>,
47 #[serde(default, skip_serializing_if = "Option::is_none")]
49 timeframe_seconds: Option<u64>,
50 #[serde(default, skip_serializing_if = "Option::is_none")]
52 tick_count: Option<u64>,
53 },
54}
55
56impl MarketEvent {
57 pub fn ts(&self) -> NaiveDateTime {
59 match self {
60 MarketEvent::Tick { ts, .. } => *ts,
61 MarketEvent::Bar { ts, .. } => *ts,
62 }
63 }
64
65 pub fn symbol(&self) -> &str {
67 match self {
68 MarketEvent::Tick { symbol, .. } => symbol,
69 MarketEvent::Bar { symbol, .. } => symbol,
70 }
71 }
72
73 pub fn to_quote(&self) -> PriceQuote {
77 self.to_quote_with_spread_fallback(None)
78 }
79
80 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 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 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 pub fn to_valid_quote(&self) -> Option<PriceQuote> {
159 let quote = self.to_quote();
160 ExecutionPricer::validate_quote("e).ok().map(|()| quote)
161 }
162}
163
164#[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 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 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 pub fn open_quote(&self) -> PriceQuote {
205 self.quote_at_mid(self.open)
206 }
207
208 pub fn close_quote(&self) -> PriceQuote {
210 self.quote_at_mid(self.close)
211 }
212}
213
214#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
238pub struct EventMetadata {
239 pub roles: SeriesRoles,
240 pub series_rank: u32,
241 pub row_sequence: u64,
243 #[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#[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 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#[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
322pub 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#[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
362pub trait DataFeed {
364 fn next_event(&mut self) -> Option<MarketEvent>;
366
367 fn peek(&self) -> Option<&MarketEvent>;
369
370 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 fn total_events(&self) -> Option<usize> {
393 None
394 }
395}
396
397pub trait FallibleBatchFeed {
399 type Error;
400
401 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#[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
430pub 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 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#[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
574pub 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 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 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#[derive(Debug, Clone)]
721pub struct VecFeed {
722 events: Vec<FeedEvent>,
723 index: usize,
724}
725
726impl VecFeed {
727 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 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 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 pub fn peek_metadata(&self) -> Option<&EventMetadata> {
757 self.events.get(self.index).map(|event| &event.metadata)
758 }
759
760 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 pub fn remaining(&self) -> usize {
780 self.events.len().saturating_sub(self.index)
781 }
782
783 pub fn total(&self) -> usize {
785 self.events.len()
786 }
787
788 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
820pub fn ticks_to_feed(ticks: Vec<data_preprocess::Tick>) -> VecFeed {
824 ticks_to_feed_with_metadata(ticks, SeriesRoles::PRIMARY, 0)
825}
826
827pub 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
856pub fn bars_to_feed(bars: Vec<data_preprocess::Bar>) -> VecFeed {
858 bars_to_feed_with_metadata(bars, SeriesRoles::PRIMARY, 0)
859}
860
861pub 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
895pub 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#[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 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, 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, 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); }
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}